import datetime import logging import time from anydi import singleton from fansifter_common.adapters.db import Database from fansifter_common.adapters.ows_account import OwsAccountClient from fansifter_common.adapters.sendgrid import SendGridClient from fansifter_common.adapters.sendgrid.exceptions import SendGridClientError from fansifter_common.adapters.sendgrid.models import Email from fansifter_common.adapters.stripo import StripoClient from fansifter_common.adapters.stripo.exceptions import StripoClientError from fansifter_common.constants import QA_ENVIRONMENT from fansifter_common.email import EmailType from fansifter_common.utils import timezone from pydantic import AwareDatetime, BaseModel from unidecode import unidecode from app.config import Settings from app.exceptions import EmailSendError, EmailSendSkipError from app.models import BatchSendRecord, CampaignBatch from app.repositories import ( BatchRecipientRepository, BatchSendRecordRepository, CampaignBatchRepository, CampaignRepository, DomainRepository, ) from app.services import FanDataService, SeedlistService from app.types import EmailSend, EmailSendResult from app.utils import batch_iterable, get_sendgrid_error_log_msg logger = logging.getLogger(__name__) class SendEmailsRequest(BaseModel): batch_id: int campaign_id: str current_hourly_quota: float emails_to_send: int dispatch_time: AwareDatetime | None = None @singleton class SendEmailsHandler: max_emails_to_send = 1000 def __init__( self, db: Database, domain_repository: DomainRepository, campaign_repository: CampaignRepository, batch_repository: CampaignBatchRepository, recipient_repository: BatchRecipientRepository, send_record_repository: BatchSendRecordRepository, sendgrid_client: SendGridClient, stripo_client: StripoClient, ows_account_client: OwsAccountClient, fandata_service: FanDataService, seedlist_service: SeedlistService, settings: Settings, ) -> None: self.db = db self.domain_repository = domain_repository self.campaign_repository = campaign_repository self.batch_repository = batch_repository self.recipient_repository = recipient_repository self.send_record_repository = send_record_repository self.sendgrid_client = sendgrid_client self.stripo_client = stripo_client self.ows_account_client = ows_account_client self.fandata_service = fandata_service self.seedlist_service = seedlist_service self.settings = settings def handle(self, request: SendEmailsRequest) -> EmailSendResult: logger.info( "Received send emails request.", extra={ "campaign_id": request.campaign_id, "batch_id": request.batch_id, "current_hourly_quota": request.current_hourly_quota, "emails_to_send": request.emails_to_send, "dispatch_time": request.dispatch_time, }, ) # Fetch batch and fan data in read session (lightweight, no locks) with self.db.session_factory(): email_send = self.create_email_send(request) # Send emails in write transaction (with DB lock for idempotency) return self.send_emails(email_send) def create_email_send(self, request: SendEmailsRequest) -> EmailSend: # noqa C901 """Create an email send for the given queue.""" batch = self.batch_repository.get(request.batch_id) if batch is None: raise EmailSendError(f"Batch {request.batch_id} is not found.") logger.info( "Creating email send for batch.", extra={ "campaign_id": batch.campaign_id, "batch_id": batch.id, "batch_size": batch.batch_size, "batch_offset": batch.batch_offset, }, ) if not batch.is_active: logger.warning( "Batch is not active, likely already processed by another invocation.", extra={ "campaign_id": batch.campaign_id, "batch_id": batch.id, "batch_size": batch.batch_size, "batch_offset": batch.batch_offset, }, ) raise EmailSendSkipError(f"Batch {batch.id} is not active.") campaign = self.campaign_repository.get(batch.campaign_id) if campaign is None: raise EmailSendError(f"Campaign {batch.campaign_id} is not found.") if not campaign.is_active: logger.warning( "Campaign %s is not active, skipping.", batch.campaign_id, extra={ "campaign_id": campaign.id, "campaign_status": campaign.status, "batch_id": batch.id, }, ) raise EmailSendSkipError(f"Campaign {campaign.id} is not active.") # Check if the email domain is found if not campaign.email_username or not campaign.sender_name: raise EmailSendError("Email username or sender details are missing.") # Check if the html and css content is missing if ( campaign.subject is None or campaign.html_content is None or campaign.css_content is None ): raise EmailSendError("Campaign subject or content is missing.") if campaign.domain_id is None: raise EmailSendError("Campaign email domain is not set.") domain = self.domain_repository.get(campaign.domain_id) if domain is None: raise EmailSendError("Campaign email domain is not found.") # Compress the html content if not compressed if not campaign.compressed_content: try: html_content = self.stripo_client.compress( html=campaign.html_content, css=campaign.css_content ) except StripoClientError as exc: raise EmailSendError( f"Campaign {campaign.id}: Failed to compress HTML content." ) from exc else: html_content = campaign.compressed_content # Get vendor data try: vendor = self.ows_account_client.get_vendor(campaign.vendor_id) except Exception as exc: raise EmailSendError( f"Campaign {campaign.id}: Failed to get vendor {campaign.vendor_id} data." ) from exc # Create the from email from_email = Email( name=campaign.sender_name, email=f"{campaign.email_username}@{domain.domain}", ) logger.info( "Retrieving fans data for batch processing", extra={ "campaign_id": batch.campaign_id, "batch_id": batch.id, "batch_size": batch.batch_size, "batch_offset": batch.batch_offset, }, ) # Get the fans data fans_data = self.fandata_service.get_fans_data( campaign_id=campaign.id, batch_id=batch.id, limit=request.emails_to_send, offset=batch.batch_offset, ) return EmailSend( campaign_id=campaign.id, batch_id=batch.id, batch_offset=batch.batch_offset, provider=batch.provider, is_first_send=batch.batch_offset == 0, brand=domain.brand, current_hourly_quota=request.current_hourly_quota, vendor_id=vendor.vendor_id, vendor_name=vendor.name, subject=campaign.subject, html_content=html_content, preview_text=campaign.preview_text, from_email=from_email, to_fans=fans_data, emails_to_send=request.emails_to_send, send_at=request.dispatch_time or timezone.now(), category=unidecode(campaign.sender_name), ) def send_emails(self, email_send: EmailSend) -> EmailSendResult: """Send emails using SendGrid.""" # Send seedlist if queue is just starting and seedlist is enabled self.send_seedlist_if_enabled(email_send) sent_emails = 0 batch_offset = email_send.emails_to_send has_failed_sends = False logger.info( "Sending email batch (%s) via SendGrid", email_send.emails_to_send, extra={ "campaign_id": email_send.campaign_id, "batch_id": email_send.batch_id, "batch_size": email_send.emails_to_send, "batch_offset": email_send.batch_offset, }, ) with self.db.transaction(): batch = self.batch_repository.get_for_update(batch_id=email_send.batch_id) if batch is None: raise EmailSendError( f"Batch {email_send.batch_id} not found during send." ) # Verify batch is still active and offset hasn't changed (idempotency check) # This prevents duplicate sends if another invocation processed this batch concurrently if not batch.is_active: logger.warning( "Idempotency check: Batch became inactive during processing, skipping.", extra={ "campaign_id": email_send.campaign_id, "batch_id": email_send.batch_id, "batch_completed_at": batch.completed_at, "batch_cancelled_at": batch.cancelled_at, "expected_offset": email_send.batch_offset, "actual_offset": batch.batch_offset, "idempotency_trigger": "batch_inactive", }, ) raise EmailSendSkipError( "Batch was already processed by another invocation." ) if batch.batch_offset != email_send.batch_offset: logger.warning( ( "Idempotency check: Batch offset changed during processing, " "skipping to prevent duplicates.", ), extra={ "campaign_id": email_send.campaign_id, "batch_id": email_send.batch_id, "expected_offset": email_send.batch_offset, "actual_offset": batch.batch_offset, "offset_drift": batch.batch_offset - email_send.batch_offset, "idempotency_trigger": "offset_changed", }, ) raise EmailSendSkipError( "Batch offset was already updated by another invocation." ) if self.should_skip_email_sending(email_send): logger.info( "Skipping real email send '%s' (automation/test mode).", email_send.subject, ) sent_emails = email_send.emails_to_send else: # Send SendGrid emails in batches iteration = 0 for to_fans, offset in batch_iterable( email_send.to_fans, batch_size=self.max_emails_to_send ): iteration += 1 try: start_time = time.perf_counter() self.sendgrid_client.send_mail( email_id=email_send.campaign_id, email_type=EmailType.CAMPAIGN, automated_email_trigger_id=None, version="v3", brand=email_send.brand, vendor_id=email_send.vendor_id, vendor_name=email_send.vendor_name, from_email=email_send.from_email, subject=email_send.subject, html_content=email_send.html_content, to_fans=to_fans, preview_text=email_send.preview_text, category=email_send.category, ) execution_time_ms = (time.perf_counter() - start_time) * 1000 logger.info( "Batch %s iteration %s: Successfully sent %s emails via SendGrid in %.2fms.", email_send.batch_id, iteration, len(to_fans), execution_time_ms, extra={ "campaign_id": email_send.campaign_id, "batch_id": email_send.batch_id, "provider": email_send.provider, "iteration": iteration, "sent_emails": len(to_fans), "execution_time_ms": round(execution_time_ms, 2), }, ) sent_emails += len(to_fans) except SendGridClientError as exc: logger.error( get_sendgrid_error_log_msg(exc), exc_info=exc, extra={ "campaign_id": email_send.campaign_id, "batch_id": email_send.batch_id, "provider": email_send.provider, "current_hourly_quota": email_send.current_hourly_quota, **exc.log_extra, }, ) batch_offset = offset has_failed_sends = True break # Update batch progress after send self.update_batch_progress( batch=batch, batch_offset=batch_offset, sent_emails=sent_emails, current_hourly_quota=email_send.current_hourly_quota, send_time=email_send.send_at, ) logger.info( "Successfully sent emails for the batch.", extra={ "campaign_id": email_send.campaign_id, "batch_id": email_send.batch_id, "batch_offset": batch_offset, "provider": email_send.provider, "sent_emails": sent_emails, "emails_to_send": email_send.emails_to_send, "current_hourly_quota": email_send.current_hourly_quota, }, ) return EmailSendResult( batch_id=email_send.batch_id, campaign_id=email_send.campaign_id, sent_emails=sent_emails, has_failed_sends=has_failed_sends, ) def update_batch_progress( self, *, batch: CampaignBatch, batch_offset: int, sent_emails: int, current_hourly_quota: float, send_time: datetime.datetime, ) -> None: """Update batch progress after sending emails.""" if batch_offset <= 0: return None last_sent_at = send_time # Move the offset batch.batch_offset += batch_offset # Update first sent at if batch.first_sent_at is None: batch.first_sent_at = send_time # Update last hour sent count if batch.last_sent_at and batch.last_sent_at.strftime( "%Y-%m-%d %H" ) == last_sent_at.strftime("%Y-%m-%d %H"): batch.last_hour_sent_emails += sent_emails else: batch.last_hour_sent_emails = sent_emails # Update last fields batch.last_sent_at = last_sent_at batch.last_hourly_quota = current_hourly_quota # Set the queue status if batch.batch_offset >= batch.batch_size: batch.completed_at = timezone.now() # Save the batch self.batch_repository.save(batch) record = BatchSendRecord( batch_id=batch.id, batch_offset=batch.batch_offset, current_hourly_quota=current_hourly_quota, sent_emails=sent_emails, sent_at=send_time, ) # Save record self.send_record_repository.save(record) return None def send_seedlist_if_enabled(self, email_send: EmailSend) -> None: """Send emails to seedlist for InboxMonster testing.""" if not (self.settings.seedlist_enabled and email_send.is_first_send): return None with self.db.session_factory(): first_batch = self.batch_repository.get_first_in_campaign( campaign_id=email_send.campaign_id ) if not (first_batch and first_batch.id == email_send.batch_id): return None # Get seedlist emails seedlist_emails = self.seedlist_service.get_seedlist_emails() if not seedlist_emails: logger.warning( "No seedlist emails found; skipping seedlist send.", extra={ "campaign_id": email_send.campaign_id, "batch_id": email_send.batch_id, }, ) return None # Create FanData for seedlist seedlist_fan_data = self.seedlist_service.create_seedlist_fan_data( seedlist_emails ) logger.info( "Sending seedlist emails.", extra={ "campaign_id": email_send.campaign_id, "batch_id": email_send.batch_id, }, ) # Send seedlist in batches for seedlist_batch, _ in batch_iterable( seedlist_fan_data, batch_size=self.max_emails_to_send ): try: self.sendgrid_client.send_mail( email_id=email_send.campaign_id, email_type=EmailType.CAMPAIGN, automated_email_trigger_id=None, version="v3", brand=email_send.brand, vendor_id=email_send.vendor_id, vendor_name=email_send.vendor_name, from_email=email_send.from_email, subject=email_send.subject, html_content=email_send.html_content, to_fans=seedlist_batch, preview_text=email_send.preview_text, category=email_send.category, is_test=True, is_seedlist=True, ) except Exception as exc: logger.error( "Failed to send seedlist emails.", exc_info=exc, extra={ "campaign_id": email_send.campaign_id, "batch_id": email_send.batch_id, }, ) break logger.info( "Seedlist emails sent successfully.", extra={ "campaign_id": email_send.campaign_id, "batch_id": email_send.batch_id, "seedlist_count": len(seedlist_emails), }, ) return None def should_skip_email_sending(self, email_send: EmailSend) -> bool: return ( self.settings.environment == QA_ENVIRONMENT and self.settings.send_skip_tag in email_send.subject )