"""Lambda function module.""" from random import random from time import sleep from common.helpers.bulk_asset import load_bulk_asset_json from common.helpers.bulk_asset import MissingProduct from common.helpers.catalog_ingestion import log_catalog_action import config from constants import jitter as jitter_const def handler(event, asset): """Lambda entrypoint.""" correlation_id = event.get('correlation_id') logger = config.get_current_logger(correlation_id) # Jitter calls jitter(correlation_id) logger.info(f'log_asset_error received: {event}') error = event.get('errors', {}).get('Error') cause = event.get('errors', {}).get('Cause') error_msg = f'{error}: {cause}' # TODO: Does this event have a product component? try: asset = load_bulk_asset_json(event) except MissingProduct: bucket = event.get('bucket') key = event.get('key') msg = f'Fatal parsing error for ' \ f'{bucket}/{key}: ' \ f'{error_msg}' logger.error(msg) lambda_result = { 'msg': error_msg, 'execution_name': event.get('execution_name'), 'state_machine_name': event.get('state_machine_name'), 'correlation_id': event.get('correlation_id') } else: log_catalog_action( config.SNOWFLAKE_S3_BUCKET, config.SNOWFLAKE_S3_LOCATION, asset, 'insert', asset.asset_type, 'Failed', error_msg ) msg = f'Fatal ingestion error for ' \ f'product_id: {asset.product.product_id} ' \ f'(upc: {asset.product.upc}). ' \ f'{error_msg}' logger.error(msg) lambda_result = { 'product_id': asset.product.product_id, 'upc': asset.product.upc, 'msg': error_msg, 'execution_name': asset.execution_name, 'state_machine_name': asset.state_machine_name, 'correlation_id': asset.correlation_id } return lambda_result def jitter(correlation_id): """Jitter lambda calls so the systems don't get hammered.""" logger = config.get_current_logger(correlation_id) if config.ENVIRONMENT.upper() == 'PROD': jitter_amt = (random() + random()) / jitter_const.PROD_FACTOR else: jitter_amt = (random() + random()) / jitter_const.QA_FACTOR logger.info(f'Jittering: {jitter_amt}') sleep(jitter_amt)