"""Start SFN execution.""" import json import time import uuid import boto3 from kafka_utils.consumer.deserializer.avro import AvroDeserializer from kafka_utils.consumer.source.mapping import EventSourceMessage from lambdacommon.common_config import logger import config def start_sfn_execution(event): """Start SFN execution.""" logger.info(f'Extracting product id from event: {event}.') product_id = _get_product_id_from_event(event) logger.info(f'Starting SFN execution for product id: {product_id}.') execution_info = _start_execution(product_id) logger.info(f'Started SFN execution with ' f'ARN: {execution_info["executionArn"]} ' f'at {execution_info["startDate"]}.') return execution_info def _get_product_id_from_event(event): """Get product id from the event.""" _, msk_message = next(iter(EventSourceMessage(event))) deserializer = AvroDeserializer(schema_registry_url=config.SCHEMA_REGISTRY_URL) event_value = deserializer.deserialize(msk_message.value, msk_message.topic, 'value') return event_value['payload']['product_id'] def _start_execution(product_id: int): """Start SFN execution.""" stepfunctions_client = boto3.client('stepfunctions', region_name=config.AWS_REGION) execution_name = f'product_id_{product_id}_timestamp_{int(time.time())}' response = stepfunctions_client.start_execution( stateMachineArn=config.STEP_FUNCTIONS_ARN, name=execution_name, input=json.dumps({'product_id': product_id}), traceHeader=f'Root=mcc-{uuid.uuid4()};Sampled=1') return response