"""Lambda send_success function module.""" import json import boto3 from lambdacommon.common_config import logger from confluent_kafka import Producer from config import BROKER_SERVERS from config import ARTIST_DIY_KAFKA_TOPIC from config import DYNAMO_TABLE dynamo_client = boto3.client('dynamodb', region_name='us-east-1') def delivery_report(err, msg): """Log delivery results from callbacks for each message produced, triggered by flush().""" if err is not None: logger.exception('Message delivery failed: {}'.format(err)) else: logger.info('Message delivered to {} [{}]'.format(msg.topic(), msg.partition())) def handler(event, context): """Lambda entry point.""" try: correlation_id = event.get('correlation_id') logger.info(f'Writing correlation id {correlation_id} to dynamodb table.') result = dynamo_client.put_item( TableName=DYNAMO_TABLE, Item={ 'correlation_id': { 'S': correlation_id } }, ReturnValues='ALL_OLD' ) if 'Attributes' in result: logger.info( 'This correlation id already exists in the gda-account-creation-success table.') logger.info('Starting to write success event to Kafka topic.') producer = Producer({'bootstrap.servers': BROKER_SERVERS, 'security.protocol': 'SSL'}) # Asynchronously produce a message, the delivery report callback # will be triggered from poll() above, or flush() below, when the message has # been successfully delivered or failed permanently. producer.produce( ARTIST_DIY_KAFKA_TOPIC, json.dumps(event).encode('utf-8'), callback=delivery_report) # Wait for any outstanding messages to be delivered and delivery report callback. producer.flush() logger.info(f'Successfully wrote success message to Kafka topic: {ARTIST_DIY_KAFKA_TOPIC}') return {'status': 'success'} except Exception as err: logger.exception(str(err)) raise err