"""Lambda module.""" from content_utils.exceptions import NoReviewQueueItemDataException from content_utils.logging.product_logging import ProductEventLogger from kafka_utils.consumer.deserializer.simple_json import JSONDeserializer from kafka_utils.consumer.source import mapping as kafka_source_mapping from kafka_utils.exceptions import IneligibleEventException import sentry_sdk from sentry_sdk.integrations.aws_lambda import AwsLambdaIntegration from src import config from src.logic import event_handling if config.SENTRY_DSN: sentry_sdk.init( dsn=config.SENTRY_DSN, environment=config.ENVIRONMENT, integrations=[AwsLambdaIntegration(timeout_warning=True)] ) deserializer = JSONDeserializer() content_lambda_logger = ProductEventLogger(config.app_logger) SKIPPED_EXCEPTIONS = ( IneligibleEventException, NoReviewQueueItemDataException, ) def handler(event, context): """Lambda entry point.""" config.app_logger.debug({'lambda_input_event': event}) msk_message = kafka_source_mapping.EventSourceMessage(event) for event_key, message in msk_message: content_lambda_logger.start(message, event_key, deserializer) try: process_event(message, deserializer) finally: content_lambda_logger.end() return {'status': 'success'} def process_event(message, deserializer): """Process MSK event.""" if not message.value: content_lambda_logger.set_data(status='skip', result='no_message_body') return try: event_handling.processing_logic(message.topic, message.value, deserializer) content_lambda_logger.set_data(status='success', result='index updated') except SKIPPED_EXCEPTIONS as e: content_lambda_logger.set_data(status='skip', result=str(e)) except Exception as e: content_lambda_logger.set_data(status='error', result=str(e)) raise e