"""Lambda trigger_gda_account_creation function module.""" import base64 from copy import deepcopy import json from datetime import datetime from json import JSONDecodeError import boto3 from confluent_kafka.cimpl import Producer from config import AWS_REGION from config import STEP_FN_ARN from config import BROKER_SERVERS from config import DEAD_LETTER_KAFKA_TOPIC from config import app_logger as logger from config import DYNAMO_TABLE CORRELATION_ID_FIELD = 'correlation_id' 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.""" initial_event = deepcopy(event) try: if not event.get('eventSource') == 'aws:kafka': logger.error('Invalid event source for this lambda. It has to be a kafka topic.') return dlq_writer = Producer({'bootstrap.servers': BROKER_SERVERS, 'security.protocol': 'SSL'}) status = {} # lambda source returns a dictionary with array of messages. # There will be 1 message in our case because of source configuration. for records in event.get('records', {}).values(): message_obj = records.pop() timestamp = datetime.fromtimestamp(message_obj.get('timestamp') / 1000) source = message_obj.get('topic') offset = message_obj.get('offset') encoded_message = message_obj.get('value') logger.info( f'Processing: offset {offset} , created at: {timestamp}. ') logger.debug(f'Message: {encoded_message}') try: record_str = base64.b64decode(encoded_message).decode('utf-8') msk_record = json.loads(record_str) logger.debug(msk_record) except (base64.binascii.Error, UnicodeDecodeError, JSONDecodeError): logger.exception(f'Failed to decode MSK message: {encoded_message}') dlq_writer.poll(0) dlq_writer.produce( DEAD_LETTER_KAFKA_TOPIC, json.dumps(event).encode('utf-8'), callback=delivery_report) status[offset] = 'failed' continue try: payload = deepcopy(msk_record) except JSONDecodeError: logger.exception(f"Failed to decode MSK payload: {msk_record.get('payload')}") dlq_writer.poll(0) dlq_writer.produce( DEAD_LETTER_KAFKA_TOPIC, json.dumps(initial_event).encode('utf-8'), callback=delivery_report) status[offset] = 'failed' continue correlation_id = payload.get(CORRELATION_ID_FIELD) dynamo_client = boto3.client('dynamodb', region_name=AWS_REGION) response = dynamo_client.get_item( TableName=DYNAMO_TABLE, Key={ 'correlation_id': { 'S': correlation_id} } ) if response.get('Item'): logger.info( f'Correlation id {correlation_id} has already successfully executed. \ Skipping Step Function Execution.') status[offset] = 'skipped' continue # if StartExecution is called with the same name and input as a running execution, the # call will succeed and return the same response as the original request. payload.update({'source': source}) # add topic name to the input kwargs = { 'stateMachineArn': STEP_FN_ARN, 'input': json.dumps(payload) } stepfn_client = boto3.client('stepfunctions', region_name=AWS_REGION) response = stepfn_client.start_execution(**kwargs) logger.info(f"State Function ARN: {response.get('executionArn')}") status[offset] = 'success' dlq_writer.flush() return {'status': status} except Exception as e: logger.exception(e) dlq_writer.poll(0) dlq_writer.produce( DEAD_LETTER_KAFKA_TOPIC, json.dumps(initial_event).encode('utf-8'), callback=delivery_report) dlq_writer.flush() return {'status': 'failed'}