import logging from typing import Any import pydantic import redis import app from app.config import settings from app.container import container from app.exceptions import EmailSendSkipError from app.handler import SendEmailsHandler, SendEmailsRequest # Set up the application app.setup() # Get the logger logger = logging.getLogger(__name__) def handle(event: Any, _context: Any) -> Any: """Handle the event.""" handler = container.resolve(SendEmailsHandler) redis_client = container.resolve(redis.Redis) try: request = SendEmailsRequest.model_validate(event) except pydantic.ValidationError as exc: logger.error("Invalid event data.", exc_info=exc, extra={"event": event}) raise lock_key = f"sender:batch:{request.batch_id}" # Lock timeout MUST be longer than Lambda timeout (2 minutes) to prevent # duplicate processing if lambda runs long. lock = redis_client.lock(lock_key, timeout=settings.redis_lock_timeout) lock_acquired = False try: lock_acquired = lock.acquire(blocking=False) if not lock_acquired: logger.warning( "Batch %s is already being processed, skipping.", request.batch_id ) return {"status": "SKIPPED"} result = handler.handle(request) if result and result.has_failed_sends: logger.error("Some emails failed to send.") raise Exception( "Email sending failed: one or more messages were not delivered." ) except EmailSendSkipError as exc: logger.info( "Skipping batch processing: %s", str(exc), extra={ "campaign_id": request.campaign_id, "batch_id": request.batch_id, }, ) return {"status": "SKIPPED"} except Exception as exc: logger.error( "Error occurred while handling the request.", exc_info=exc, extra={ "campaign_id": request.campaign_id, "batch_id": request.batch_id, }, ) raise finally: if lock_acquired: if lock.owned(): lock.release() else: logger.warning("Lock %s already expired or not owned.", lock_key) return {"status": "OK"}