"""Main application logic for processing triggered sends from S3 events.""" from collections.abc import Iterator import json import logging from typing import Optional import boto3 from lambdacommon.common_config import logger import sentry_sdk from sentry_sdk.integrations.aws_lambda import AwsLambdaIntegration from sentry_sdk.integrations.logging import LoggingIntegration import config # noqa from src.kafka_producer import TriggeredSendEventProducer from src.marketing_cloud_api import triggered_sends from src.models import TriggeredSend from src.models import TriggeredSendFailedEvent from src.models import TriggeredSendSuccessfulEvent logging_integration = LoggingIntegration( level=logging.INFO, # Capture info and above as breadcrumbs event_level=logging.CRITICAL # Send only critical as events ) sentry_sdk.init( integrations=[AwsLambdaIntegration(), logging_integration] ) def handler(event, context): """Lambda entry point. The Lambda is triggered by an S3 event. It reads the uploaded file and processes its contents to initiate triggered sends through the Marketing Cloud API, based on the data in the file. """ error = None # Check for 'Records' in event if not event.get('Records'): error = "Event does not contain 'Records' key or it is empty." logger.error(f'Error: {error}') return { 'statusCode': 400, 'body': json.dumps({'message': error}) } s3_key = event['Records'][0] bucket_name = s3_key['s3']['bucket']['name'] object_key = s3_key['s3']['object']['key'] logger.info(f'Processing file {object_key} from bucket {bucket_name}') try: triggered_send_info_records = read_s3_file(bucket_name, object_key) process_triggered_sends(triggered_send_info_records) except Exception as e: error = str(e) logger.critical(f'Error: {error}') finally: # Always delete the file after processing delete_s3_file(bucket_name, object_key) message = f'File {object_key} from bucket {bucket_name} processed successfully.' if error: message = f'File {object_key} from bucket {bucket_name} failed to process with error: {error}' return { 'statusCode': 200, # Failure is handled separately 'body': json.dumps({'message': message}) } def process_triggered_sends(triggered_send_info_records: Iterator[str]) -> None: """Process triggered sends based on the records read from S3. This function iterates through the records, sending each one to the Marketing Cloud API. It also produces events to Kafka for successful and failed sends. Args: triggered_send_info_records (Iterator[str]): A list of triggered send information. """ logger.debug('Initializing Kafka producer for triggered sends...') kafka_producer = TriggeredSendEventProducer() for i, triggered_send_info_raw in enumerate(triggered_send_info_records, start=1): logger.debug(f'Loading triggered send record #{i}') triggered_send_info = TriggeredSend.model_validate_json(triggered_send_info_raw) _process_triggered_send(triggered_send_info, kafka_producer) def _process_triggered_send(triggered_send_info: TriggeredSend, kafka_producer: TriggeredSendEventProducer) -> None: """Process a single triggered send and produce Kafka events based on the result. Args: triggered_send_info (TriggeredSend): The information for the triggered send. kafka_producer (TriggeredSendEventProducer): The Kafka producer instance. """ logger.debug('Submitting triggered send.') submission_response, submission_error = triggered_sends.submit_triggered_send(triggered_send_info) if submission_error: _handle_failed_send(triggered_send_info, submission_error, kafka_producer) else: submission_response_data, validation_error = ( triggered_sends.triggered_send_response_validation(submission_response)) if validation_error: _handle_failed_send(triggered_send_info, validation_error, kafka_producer, submission_response_data) else: _handle_successful_send(triggered_send_info, kafka_producer) def _handle_failed_send( triggered_send_info: TriggeredSend, error: str, kafka_producer: TriggeredSendEventProducer, api_response: Optional[dict] = None ) -> None: """Handle a failed triggered send and produce a Kafka failed event. Args: triggered_send_info (TriggeredSend): The information for the triggered send. error (str): The error encountered during the send. kafka_producer (TriggeredSendEventProducer): The Kafka producer instance. api_response (Optional[dict]): Optional API response payload. """ logger.error(f'Triggered send failed with error: {error}') failed_event = TriggeredSendFailedEvent( triggered_send_info=triggered_send_info, error=error, api_response=api_response ) kafka_producer.produce_failed_event(failed_event) def _handle_successful_send(triggered_send_info: TriggeredSend, kafka_producer: TriggeredSendEventProducer) -> None: """Handle a successful triggered send and produce a Kafka successful event. Args: triggered_send_info (TriggeredSend): The information for the triggered send. kafka_producer (TriggeredSendEventProducer): The Kafka producer instance. """ logger.info('Triggered send was successful.') successful_event = TriggeredSendSuccessfulEvent(triggered_send_info=triggered_send_info) kafka_producer.produce_successful_event(successful_event) def read_s3_file(bucket_name: str, object_key: str) -> Iterator[str]: """Read a file from S3 and yield its content line by line. Args: bucket_name (str): The name of the S3 bucket. object_key (str): The key of the object in the S3 bucket. Yields: str: Each line of the file as a decoded string. """ logger.debug(f'Reading file {object_key} from bucket {bucket_name}...') s3_client = boto3.client('s3') response = s3_client.get_object( Bucket=bucket_name, Key=object_key, ExpectedBucketOwner=config.AWS_ACCOUNT_ID ) for i, line in enumerate(response['Body'].iter_lines()): yield line.decode('utf-8') def delete_s3_file(bucket_name: str, object_key: str): """Delete a file from S3. Args: bucket_name (str): The name of the S3 bucket. object_key (str): The key of the object in the S3 bucket. """ s3_client = boto3.client('s3') s3_client.delete_object( Bucket=bucket_name, Key=object_key, ExpectedBucketOwner=config.AWS_ACCOUNT_ID )