"""Lambda review-queue-populate function module.""" import json import uuid from content_utils.logging.product_logging import ProductEventLogger from kafka_utils.consumer.deserializer.avro import AvroDeserializer from kafka_utils.consumer.source.mapping import EventSourceMessage import sentry_sdk from sentry_sdk.integrations.aws_lambda import AwsLambdaIntegration from src import config from src import constants from src.logic import event_processing if config.SENTRY_DSN: sentry_sdk.init( dsn=config.SENTRY_DSN, environment=config.ENVIRONMENT, integrations=[AwsLambdaIntegration(timeout_warning=True)] ) content_lambda_logger = ProductEventLogger(config.app_logger) avro_deserializer = AvroDeserializer(config.SCHEMA_REGISTRY_URL) def handler(message, context): """Lambda Entry point.""" correlation_id = str(uuid.uuid4()) msk_message = EventSourceMessage(message) error_outputs = [] for key, event in msk_message: content_lambda_logger.start(event, key, avro_deserializer) result = process_event(event, correlation_id) status = result.get('status') content_lambda_logger.set_data( status=status, product_id=result.get('product_id'), queue_id=result.get('review_queue_id'), result=result.get('message') ) if status == constants.STATUS_ERROR: error_outputs.append(result) content_lambda_logger.end() if len(error_outputs): raise Exception(json.dumps(error_outputs)) def process_event(event, correlation_id): """Process Kafka event.""" if not event.value: return { 'status': constants.STATUS_SKIP, 'message': 'no_message_body' } return event_processing.process_event( event.topic, event.value, avro_deserializer, correlation_id )