"""Lambda function module.""" import json import os from random import random from time import sleep from types import SimpleNamespace from common.helpers.asset import copy_asset from common.helpers.asset import create_asset_token from common.helpers.bulk_asset import load_bulk_asset_json from common.lambda_exceptions import AudioException import config from config import graphql_gateway from constants import jitter as jitter_const from lambdacommon.graphql.graphql import GraphQLError def handler(event, context): """Lambda entrypoint.""" correlation_id = event.get('correlation_id') # DDEX_INGESTER_INTEGRATION: Implemented standard logger, not # DDEXIngesterAdapter logger logger = config.get_current_logger(correlation_id) # Jitter calls jitter(correlation_id) graphql_gateway.set_headers( { 'Orchard-User-Id': config.OA_USER, 'Correlation-Id': correlation_id } ) # Log lambda begins message logger.info(f'handle_audio received: {event}') return handle_track_asset(event, correlation_id) def handle_track_asset(event, correlation_id=None): """Handle uploading track audio asset.""" logger = config.get_current_logger(correlation_id) asset = load_bulk_asset_json(event) track = asset.product.track if not track.tuid: logger.error(f'TUID not found for ISRC: {track.isrc}') raise AudioException( f'TUID not found for ISRC: {track.isrc}' ) logger.info(f'Asset source bucket: {track.bucket}') logger.info(f'Asset key: {track.key}') # Ask for a token from ows-assets try: asset_token = create_asset_token(graphql_gateway, 'audio') except GraphQLError as e: error_str = 'GraphQL create asset error' msg = f'{error_str}: {str(e)}' logger.error(msg) raise e filename = asset_token['filename'] + os.path.splitext(track.filename)[1] target_bucket = asset_token['bucket'] # Construct metadata and copy our asset to target bucket from ows-assets s3_metadata = construct_s3_metadata( asset.product, track ) try: copy_asset(s3_metadata, track, target_bucket, filename) except Exception as e: error_str = 'S3 Copy error' msg = f'{error_str}: {str(e)}' logger.error(msg) raise e track.ows_assets_filename = filename asset.product.track = track return json.loads( json.dumps(asset, default=lambda s: vars(s))) # DDEX_INGESTER_INTEGRATION: Modified to remove is_correction param def construct_s3_metadata( product, track) -> dict: """Construct metadata dict for s3.""" if not track.filename.isascii(): track_filename = track.filename.encode( 'ascii', 'replace').decode('ascii') else: track_filename = track.filename return { 'asset_type': format_asset_type(track), 'product_id': str(product.product_id), 'upc': str(product.upc), 'track_unique_id': track.tuid, 'original_filename': track_filename, 'is_correction': '0', } def format_asset_type(asset: SimpleNamespace): """Format asset type from asset object. Args: asset (object): Asset object Returns: dict """ return os.path.splitext(asset.filename)[1].strip('.').upper() 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)