"""process-participants.""" import json import uuid import boto3 import config from config import graphql_gateway from constants import queries from ddex_ingester_common.constants.ddex_providers import SME from ddex_ingester_common.helpers.s3_ddex import load_ddex_json from ddex_ingester_common.lambda_exceptions import ProcessParticipantsException from ddex_ingester_common.logging import utils as logging_utils from ddex_ingester_common.schemas.s3_schema import S3Schema from ddex_ingester_common.schemas.state_machine_schema import ( StateMachineSchema ) from lambdacommon.graphql import graphql logger = logging_utils.get_logger(config.app_logger) def handler(event, context): """Lambda entry point.""" parsed_ddex = S3Schema().load(load_ddex_json(event)) context = StateMachineSchema().load(event) correlation_id = context.correlation_id or str(uuid.uuid4()) context.correlation_id = correlation_id logging_utils.update_logger_correlation_id(logger, correlation_id) logging_utils.update_logger_with_message_ids( logger, parsed_ddex.message_id, parsed_ddex.message_thread_id, parsed_ddex.execution_name ) graphql_gateway.set_headers( { 'Orchard-User-Id': config.OA_USER, 'Correlation-Id': correlation_id, } ) participants = get_participants(parsed_ddex) participants_dict = {} participants_renaming = {} vendor_id = context.product.vendor_id subaccount_id = context.product.subaccount_id try: for participant in participants: logger.info(f'Processing Participant: {participant}') artist = get_artist(participant, vendor_id) if not artist: artist = create_artist( participant, vendor_id, subaccount_id ) label_participant = create_label_participant( participant, vendor_id, subaccount_id ) if label_participant.get('id') not in participants_dict: participants_dict[label_participant.get('id')] = { 'name': participant['name'], 'artist_id': artist.get('artistId'), 'label_participant_id': label_participant.get('id'), 'label_participant_uuid': label_participant.get('uuid'), } else: if participants_dict[label_participant.get('id')]['name'] != \ participant['name']: participants_renaming[participant['name']] = \ participants_dict[label_participant.get('id')]['name'] except graphql.GraphQLError as err: if context.ddex_provider == SME: raise ProcessParticipantsException('Graphql error') from err raise except Exception as exp: if context.ddex_provider == SME: raise ProcessParticipantsException( f'Error processing participants.\n{exp}') from exp raise parsed_ddex.label_participants = [p for p in participants_dict.values()] rename_participants(parsed_ddex, participants_renaming) remove_duplicates(parsed_ddex) context.correlation_id = correlation_id save_s3_context(context, parsed_ddex) return StateMachineSchema().dump(context) def get_participants(parsed_ddex): """Merge all necessary data for each participant into a dict. Args: parsed_ddex(S3Schema): deserialized ddex from S3 Returns: list """ participants_dict = {} all_participants = get_all_participants(parsed_ddex) for participant in all_participants: name = participant.name if participants_dict.get(name) is None: participants_dict[name] = { 'name': name, 'roles': participant.roles, } else: check_participant_ids(participant, participants_dict[name]) if participant.apple_id: participants_dict[name]['apple_id'] = participant.apple_id if participant.spotify_uri: participants_dict[name]['spotify_uri'] = participant.spotify_uri return [p for p in participants_dict.values()] def get_all_participants(parsed_ddex): """Collect all participants from parsed_ddex, including duplicates. Args: parsed_ddex(S3Schema): deserialized ddex from S3 Returns: list """ all_participants = [] if parsed_ddex.project and parsed_ddex.project.artist: all_participants.append(parsed_ddex.project.artist) if parsed_ddex.product.display_artists: all_participants.extend(parsed_ddex.product.display_artists) for track in parsed_ddex.tracks: if track.display_artists: all_participants.extend(track.display_artists) if track.resource_contributors: all_participants.extend(track.resource_contributors) if parsed_ddex.video: if parsed_ddex.video.display_artists: all_participants.extend(parsed_ddex.video.display_artists) if parsed_ddex.video.resource_contributors: all_participants.extend(parsed_ddex.video.resource_contributors) return all_participants def check_participant_ids(participant, participant_in_dict): """Check for different apple_ids or spotify_uris and raise an Exception. Args: participant (S3Schema Participant): Participant to check ids participant_in_dict (dict): Current information about participant """ def get_roles_str(roles): return ', '.join(f'"{i}"' for i in roles or []) or '""' name = participant.name apple_id_1 = participant.apple_id apple_id_2 = participant_in_dict.get('apple_id') if apple_id_1 and apple_id_2 and apple_id_1 != apple_id_2: roles_1 = get_roles_str(participant.roles) roles_2 = get_roles_str(participant_in_dict['roles']) message = ( f'Multiple Apple IDs for participant "{name}": ' f'"{apple_id_1}" with roles {roles_1} and ' f'"{apple_id_2}" with roles {roles_2}' ) logger.error(message) raise ProcessParticipantsException(message) spotify_uri_1 = participant.spotify_uri spotify_uri_2 = participant_in_dict.get('spotify_uri') if ( spotify_uri_1 and spotify_uri_2 and spotify_uri_1 != spotify_uri_2 ): roles_1 = get_roles_str(participant.roles) roles_2 = get_roles_str(participant_in_dict['roles']) message = ( f'Multiple Spotify URIs for participant "{name}": ' f'"{spotify_uri_1}" with roles {roles_1} and ' f'"{spotify_uri_2}" with roles {roles_2}' ) logger.error(message) raise ProcessParticipantsException(message) def get_artist(participant, vendor_id): """Format payload and call GraphQL query to get artist data. Args: participant (dict): dict with participant data vendor_id (str): Product vendor_id Returns: dict """ payload = { 'artistName': participant.get('name'), 'vendorId': vendor_id } logger.info(f'Executing filterArtists with payload: {payload}') result = graphql_gateway.execute(queries.get_artist, payload) logger.info(f'filterArtists returned: {result}') if result['data']['filterArtists']: return result['data']['filterArtists'][0] else: return None def create_artist( participant, vendor_id, subaccount_id): """Format payload and call GraphQL mutation to create artist. Args: participant (dict): dict with participant data vendor_id (str): Product vendor_id subaccount_id (str): Product subaccount_id Returns: dict """ payload = { 'create': [ { 'artistName': participant.get('name'), 'vendorId': vendor_id, 'subaccountId': subaccount_id, } ] } logger.info(f'Executing saveArtists with payload: {payload}') result = graphql_gateway.execute(queries.save_artist, {'data': payload}) logger.info(f'saveArtists returned: {result}') if result['data']['saveArtists']: return result['data']['saveArtists'][0] else: return {} def create_label_participant( participant, vendor_id, subaccount_id): """Format payload and call GraphQL mutation to create label participant. Args: participant (dict): dict with participant data vendor_id (str): Product vendor_id subaccount_id (str): Product subaccount_id Returns: dict """ payload = { 'name': participant.get('name'), 'spotifyId': participant.get('spotify_uri'), 'appleMusicId': participant.get('apple_id') } # Remove keys with null values payload = { key: value for key, value in payload.items() if value is not None} 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 rename_participants(s3_context, participants_renaming): """Rename participants to be consistent with label_participant_id.""" participants = get_all_participants(s3_context) for participant in participants: if participant.name in participants_renaming: participant.name = participants_renaming[participant.name] def remove_duplicates(s3_context): """Remove duplicated artists after renaming.""" def _remove_duplicates(artists): deduplicated_artists = [] added_artists = [] for artist in artists: if artist.name not in added_artists: added_artists.append(artist.name) deduplicated_artists.append(artist) return deduplicated_artists if s3_context.product.display_artists: s3_context.product.display_artists = _remove_duplicates( s3_context.product.display_artists) for track in s3_context.tracks: if track.display_artists: track.display_artists = _remove_duplicates( track.display_artists) def save_s3_context(context, s3_context): """Save s3 context.""" s3_client = boto3.client('s3') json_file_path = f'{context.key}parsed_ddex.json' s3_client.put_object( Bucket=context.bucket, Key=json_file_path, Body=json.dumps(S3Schema().dump(s3_context)).encode(encoding='UTF-8') )