"""Lambda function module.""" from random import random from time import sleep from types import SimpleNamespace from typing import Dict from typing import Optional from common.helpers.asset import get_product_track_assets from common.helpers.bulk_asset import load_bulk_asset_json from common.helpers.catalog_ingestion import log_catalog_action from common.lambda_exceptions import AudioException import config from config import graphql_gateway from constants import jitter as jitter_const from constants.asset import ENCODING_COMPLETED def handler(event, context): """Lambda entrypoint.""" correlation_id = event['correlation_id'] logger = config.get_current_logger(correlation_id) logger.info(f'poll_audio_status received: {event}') # Jitter calls jitter(correlation_id) asset = load_bulk_asset_json(event) graphql_gateway.set_headers( { 'Orchard-User-Id': config.OA_USER, 'Correlation-Id': correlation_id } ) track = asset.product.track # Call ows assets for status of uploaded track asset status = check_audio_valid(asset.product, track.tuid, correlation_id) # Log success condition success_msg = f'Audio asset for Product Id {asset.product.product_id} ' \ f'UPC {asset.product.upc} ' \ f'TUID {asset.product.track.tuid} ' \ f'Track Key {asset.product.upc}_' \ f'{asset.product.track.track_volume_number}_' \ f'{asset.product.track.track_sequence_number} ' \ f'transferred successfully: {status}' logger.info(success_msg) # No error - Write success json log_catalog_action( config.SNOWFLAKE_S3_BUCKET, config.SNOWFLAKE_S3_LOCATION, asset, 'insert', 'audio', 'Success', success_msg ) lambda_result = { 'product_id': asset.product.product_id, 'upc': asset.product.upc, 'tuid': asset.product.track.tuid, 'msg': success_msg, 'execution_name': asset.execution_name, 'state_machine_name': asset.state_machine_name, 'correlation_id': asset.correlation_id } return lambda_result # DDEX_INGESTER_INTEGRATION: Replaces get_asset_status() def check_audio_valid(product: SimpleNamespace, tuid: int, correlation_id: Optional[str] = None) -> Dict: """Get track asset status.""" logger = config.get_current_logger(correlation_id) # Default to V2 use_v2 = True logger.info('V2 audio polling is enabled') return poll_asset_status( use_v2, product.product_id, tuid, correlation_id) def poll_asset_status(use_v2: bool, product_id: int, tuid: int, correlation_id: Optional[str] = None) -> Dict: """V1 flow for asset status checking.""" finished_status = ENCODING_COMPLETED logger = config.get_current_logger(correlation_id) logger.info("Checking product's assets") assets = get_product_track_assets( graphql_gateway, str(product_id), use_v2 ) if not assets: raise AudioException('Assets not yet available!') logger.info(f'Audio assets: {assets}') status = next( (asset for asset in assets if asset.get('trackUniqueId') == int(tuid)), None ) if not status: raise AudioException('Assets not yet available!') # Note: V2 endpoint does not have a correction field. # Replacement asset replaces the original in the GraphQL response. logger.info(f'Audio asset status: {status}') if status.get('status') != finished_status: raise AudioException(f'Asset status is not {finished_status}') return status 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)