"""lookup-orchard-product.""" import json from os import path from random import random from time import sleep import uuid import boto3 from common.constants import fields from common.constants import queries from common.lambda_exceptions import SetTrackMetadataException import config from config import graphql_gateway from constants import jitter as jitter_const from helpers.bulk_asset import load_graphql_product_result from lambdacommon.graphql.graphql import GraphQLError def handler(event, context): """Lambda entry point.""" bucket = event.get('bucket') key = event.get('key') item_index = event.get('item_index') execution_name = event.get('execution_name') state_machine_name = event.get('state_machine_name') correlation_id = event.get('correlation_id') or str(uuid.uuid4()) graphql_gateway.set_headers( { 'Orchard-User-Id': config.OA_USER, 'Correlation-Id': correlation_id } ) # DDEX_INGESTER_INTEGRATION: Implemented standard logger, not # DDEXIngesterAdapter logger logger = config.get_current_logger(correlation_id) logger.info(f'lookup_product received {bucket}/{key}.') # Jitter calls jitter(correlation_id) s3_client = boto3.client('s3') # TODO: log Catalog Ingestion Error if S3 cannot find file? file = s3_client.get_object(Bucket=bucket, Key=key) file_content = file.get('Body').read() data = json.loads(file_content) item = map_keys(data[item_index]) upc = item.get('upc') try: graphql_result = get_product_metadata(upc, correlation_id) except GraphQLError as e: error_str = 'GraphQL get product error' msg = f'{error_str}: {str(e)}' logger.error(msg) raise e # TODO: create fallback logic for product not found. # - Check if product_id has been passed. Try with product_id if not graphql_result or not graphql_result.get('productId'): error_str = 'Product not found' msg = f'{error_str}: {upc}' logger.error(msg) raise ValueError(f'Product for {upc} not found.') # DDEX_INGESTER_INTEGRATION: Skip ddex-ingester context JSON, create basic # JSON element structure product = load_graphql_product_result(graphql_result) asset = dict() asset['bucket'] = item.get('bucket') asset['key'] = path.join( item.get('path'), item.get('filename')) asset['filename'] = item.get('filename') asset['ows_assets_filename'] = None if item['asset_type'] == 'cover': product['artwork'] = asset elif item['asset_type'] == 'audio': product['track'] = asset track = get_track_metadata(graphql_result, item) product['track']['isrc'] = track.get('isrc') product['track']['tuid'] = track.get('tuid') product['track']['track_sequence_number'] = track.get('trackNumber') product['track']['track_volume_number'] = track.get('volumeNumber') product['track']['track_name'] = track.get('trackName') product_id = product['product_id'] asset_type = item['asset_type'] logger.info( f'lookup_product found: ' f'Asset type {asset_type}, ' f'UPC {upc}, ' f'Product Id: {product_id}.') return { 'product': product, 'asset_type': item['asset_type'], 'execution_name': execution_name, 'state_machine_name': state_machine_name, 'correlation_id': correlation_id } def get_product_metadata(upc, correlation_id=None): """Get product metadata via GraphQL. Args: upc (str): UPC, correlation_id (str): a correlation_id UUID4 Returns: result (dict): GraphQL response """ # DDEX_INGESTER_INTEGRATION: Implemented standard logger, not # DDEXIngesterAdapter logger logger = config.get_current_logger(correlation_id) result = graphql_gateway.execute( queries.get_product_by_upc, {'upc': str(upc)} )['data']['productByUpc'] logger.info(f'Graphql result: {result}') return result # DDEX_INGESTER_INTEGRATION: Removed unneeded fields def map_keys(data): """Map keys to proper case from input data. Args: data (dict): Input data Returns: dict """ return { 'upc': data.get(fields.upc_field), 'filename': data.get(fields.file_name), 'bucket': data.get(fields.s3_source_bucket), 'path': data.get(fields.s3_source_file_path), 'asset_type': data.get(fields.asset_type), 'file_size': data.get(fields.file_size), 'foreign_key': data.get(fields.foreign_rel_id), 'product_id': data.get(fields.product_id), 'project_code': data.get(fields.project_code), 'session_id': data.get(fields.session_id), 'seq_num': data.get(fields.seq_num), 'upload_date': data.get(fields.upload_date), 'vendor_id': data.get(fields.vendor_id), 'volume_num': data.get(fields.volume_num) } 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) def get_track_metadata(graphql_result, item): """Parse GraphQL result for track metadata. Args: graphql_result (dict): A graphql get_product result item (dict): The item being looked up Returns: dict The appropriate track metadata from the list of tracks. """ track_list = graphql_result.get('tracks', []) if not len(track_list): upc = graphql_result.get('upc') raise SetTrackMetadataException(f'Track list not found for {upc}') for track in track_list: track_match = track['trackNumber'] == item['seq_num'] vol_match = track['volumeNumber'] == item['volume_num'] if track_match and vol_match: return track # Failed to fin match track_num = item['seq_num'] vol_num = item['volume_num'] upc = item['upc'] raise SetTrackMetadataException( f'No match for volume {vol_num}, ' f'track {track_num} ' f'on upc {upc}')