"""processs-track-performers.""" import json import os from typing import Dict import uuid import boto3 from pandas import isna from lambdacommon.graphql import graphql from constants.file import ( INPUT_FILE_EXT_LIST ) from constants import fields from constants.role_mappings import ( BULK_PERFORMER_ROLE_TO_ORCHARD_PERFORMER, BULK_PERFORMER_TYPE_TO_ORCHARD_PERFORMER_TYPE, DDEX_RESOURCE_CONTRIBUTOR_ROLE_TO_ORCHARD_PERFORMER_ROLE, PRIMARY_ARTIST_ROLE_MAP, WRITER_ROLE_MAP ) from parse_xlsx import parse_perfomer_metadata import logging_utils from constants import queries from config import graphql_gateway import config logger = logging_utils.get_logger(config.app_logger) class ProcessParticipantsException(Exception): """ProcessParticipants lambda exception.""" def main(): """Parse XLSX performer update.""" correlation_id = config.DEV_CORRELATION_ID or str(uuid.uuid4()) logging_utils.update_logger_correlation_id(logger, correlation_id) graphql_gateway.set_headers( { 'Orchard-User-Id': config.OA_USER, 'Correlation-Id': correlation_id, "apollographql-client-name": "lvona_test" } ) logger.info(f'Processing performer updates from {config.KEY}') bucket = config.BUCKET key = config.KEY vendor_id = config.VENDOR_ID subaccount_id = config.SUBACCOUNT_ID check_key_extension(key) s3_client = boto3.client('s3') xlsx_data = s3_client.get_object(Bucket=bucket, Key=key).get('Body').read() # convert XLSX to dict via pandas performer_dict = parse_perfomer_metadata(xlsx_data) track_performers = get_tracks_with_performers(performer_dict) track_update_list = get_track_participants(track_performers, vendor_id) process_track_updates(track_update_list) return def check_key_extension(key): """Check file extension is a valid type.""" _, extension = os.path.splitext(key) if extension not in INPUT_FILE_EXT_LIST: msg = 'S3 metadata file is not a valid file type' logger.warning(msg) raise TypeError(msg) def get_tracks_with_performers(performer_data): """Merge all necessary data for each track/participant pair into a dict. Args: performer_data(list): data from XLSX Returns: dict: track/performer data keyed off upc """ product_dict = {} track_list = {} for track in performer_data: # get current UPC upc = track.get(fields.DIGITAL_UPC) isrc = track.get(fields.ISRC) # Add UPC section if missing if not product_dict.get(upc): product_dict[upc] = get_product_by_upc(str(upc))['tracks'] tuid = [i['tuid'] for i in product_dict.get(upc) if i['isrc'] == isrc][0] # noqa if not track_list.get(tuid): track_list[tuid] = [] # Create unique item for each performer on a single track row for i in range(1, fields.PERFORMER_FIELD_COUNT+1): name_field=f'{fields.PERFORMER_FIELD_PREFIX} {i} ' \ f'{fields.PERFORMER_FIELD_SUFFIX_NAME}' type_field = f'{fields.PERFORMER_FIELD_PREFIX} {i} ' \ f'{fields.PERFORMER_FIELD_SUFFIX_TYPE}' role_field = f'{fields.PERFORMER_FIELD_PREFIX} {i} ' \ f'{fields.PERFORMER_FIELD_SUFFIX_ROLE}' name = track.get(name_field) perf_type = track.get(type_field) role = track.get(role_field) # Cull blank artists if not isna(name) and not isna(perf_type) and not isna(role): track_list[tuid].append( { 'name': name, 'type': perf_type, 'role': role } ) return track_list def get_track_participants(track_performers, vendor_id): """Get the participants for a list of tracks. Args: track_performers (dict): Track/performer data keyed off upc. vendor_id (int): The vendor id associated with the participants. Returns: dict: Track/performer data with participants attached. """ track_update_list = {} # Sanity - segment by vendor to prevent collisions label_participant_dict = { vendor_id: {} } for tuid, performers in track_performers.items(): if not track_update_list.get(tuid): track_update_list[tuid] = { 'participant_data': [], 'performer_data': [] } # Get track metadata for participants track_data = get_track_data(tuid) # Handle Performers for performer in performers: performer_name = performer.get('name') track_update_list[tuid]['performer_data'].append( { 'name': performer_name, 'roleId': DDEX_RESOURCE_CONTRIBUTOR_ROLE_TO_ORCHARD_PERFORMER_ROLE.get( BULK_PERFORMER_ROLE_TO_ORCHARD_PERFORMER.get( performer.get('role'))).get('roleId'), 'type': BULK_PERFORMER_TYPE_TO_ORCHARD_PERFORMER_TYPE.get(performer.get('type')) } ) # Handle Primary Artists for artist in track_data.get('primaryArtists'): artist_name = artist.get('artistName') # Check if label_participant has been checked already if not label_participant_dict[vendor_id].get(artist_name): label_participant = create_label_participant( artist_name, vendor_id ) # Store the label_participant info label_participant_dict[vendor_id][artist_name] = \ label_participant # Convert `primaryArtists` field info to track update mutation syntax track_update_list[tuid]['participant_data'].append( { 'labelParticipantUuid': label_participant_dict[vendor_id][artist_name]['uuid'], 'role': PRIMARY_ARTIST_ROLE_MAP[artist.get('artistType')] } ) # Handle Writers for writer in track_data.get('writers'): writer_name = writer.get('name') # Check if label_participant has been checked already if not label_participant_dict[vendor_id].get(writer_name): label_participant = create_label_participant( writer_name, vendor_id ) # Store the label_participant info label_participant_dict[vendor_id][writer_name] = \ label_participant # Convert `writers` field info to track update mutation syntax track_update_list[tuid]['participant_data'].append( { 'labelParticipantUuid': label_participant_dict[vendor_id][writer_name]['uuid'], 'role': WRITER_ROLE_MAP[writer.get('type')] } ) return track_update_list def get_product_by_upc(upc): """Format payload and call GraphQL query to get product data by upc. Args: upc (int): UPC Returns: dict """ payload = { 'upc': upc } logger.info(f'Executing productByUpc with payload: {payload}') result = graphql_gateway.execute(queries.GET_PRODUCT_BY_UPC, payload) logger.info(f'productByUpc returned: {result}') if result['data']['productByUpc']: return result['data']['productByUpc'] else: return None def create_label_participant( performer_name, vendor_id, subaccount_id = None): """Format payload and call GraphQL mutation to create label participant. Args: performer_name (str): Name of participant. vendor_id (str): Product vendor_id subaccount_id (str): Product subaccount_id Returns: dict """ payload = { 'name': performer_name } logger.info(f'Executing createLabelParticipant with payload: {payload}') result = graphql_gateway.execute( queries.GET_OR_CREATE_LABEL_PARTICIPANT, { 'data': payload, 'vendorId': vendor_id, 'subaccountId': subaccount_id or 0, # Mutation expects 0 if there is no subaccountId } ) logger.info(f'createLabelParticipant returned: {result}') return result['data']['createLabelParticipant'] def get_track_data(tuid: int) -> Dict: """Retrieve the data of a track.""" result = graphql_gateway.execute( queries.GET_TRACK_BY_TUID, {'tuid': tuid} )['data']['track'] logger.info(f'Ran get track with tuid {tuid} and received {result}') return result def update_track_metadata(payload): """Update the metadata of a given track.""" logger.info(f'Update track metadata GraphQL mutation: {payload}') graphql_gateway.execute( queries.UPDATE_TRACK_METADATA, {'data': payload} ) def process_track_updates(track_update_list: Dict): """Update tracks with new performers. Raises: ProcessParticipantsException: GraphQL Error """ for tuid, artist_data in track_update_list.items(): track_update = { 'update': { 'tracks': [tuid], 'body': { 'participations': artist_data.get('participant_data'), 'performers': artist_data.get('performer_data') } } } try: update_track_metadata(track_update) except graphql.GraphQLError as err: err_msg = f'GraphQL Could not process updates for tuid: {tuid}' logger.error(err_msg) raise ProcessParticipantsException('Graphql error') from err except Exception as exp: raise ProcessParticipantsException( f'Error processing participants.\n{exp}') from exp def delete_object(bucket: str, key: str) -> object: """Delete S3 object by filename (key).""" s3_client = boto3.client('s3') response = s3_client.delete_object( Bucket=bucket, Key=key ) return response if __name__ == '__main__': from datetime import datetime startTime = datetime.now() main() print(f'Total time: {datetime.now() - startTime}')