"""Lambda function to get GRPS data based on UPC.""" import uuid from common.connectors.snowfalke_connector import execute_snowflake_query from common.models.state_machine.grps_ingestion_context import \ GrpsIngestionContext from common.models.state_machine.participant import Participant from common.models.state_machine.product import Product from common.models.state_machine.project import Project from common.models.state_machine.track import Track from common.schemas.state_machine_schema import ProductType, StateMachineSchema import config from src.constants.grps_queries import \ GET_ART_RELATIONS_PRODUCT_BY_PRODUCT_CODE, GET_LABEL_BY_UPC, \ GET_PARENT_PARENT_AND_REP_OWNER_KEY_BY_UPC, GET_PRODUCT_DATA_BY_UPC, \ GET_PRODUCT_PARTICIPANTS_BY_UPC, GET_PRODUCT_TRACKS_BY_UPC, \ GET_PRODUCT_TRACKS_PARTICIPANTS_BY_UPC_ISRC, GET_PROJECT_BY_PRODUCT_UPC, \ GET_PROJECT_PARTICIPANT_BY_REC_PROJECT_ID, \ GET_VENDOR_AND_SUBACCOUNT_BY_PARENT_AND_REP_OWNER_KEY from src.exceptions import AmbiguousVendorOrSubaccountMapping, \ GetGrpsDataException, RepOwnerDoNotIngestException from src.prep_grps import prep_forbidden_sequence, prep_track_name_length, \ prep_track_sequence_and_volume_number, remove_track_isrc_duplicates logger = config.app_logger def handler(event, input_context): """Lambda function to get GRPS data based on UPC.""" correlation_id = str(uuid.uuid4()) upc = event['upc'] grps_ingestion_id = event['grps_ingestion_id'] prod_no = event['prod_no'] product, product_type = get_product(upc, prod_no) project = get_project(upc, product, prod_no) tracks = get_tracks(upc, prod_no, product.release_type) tracks = remove_track_isrc_duplicates(tracks) prep_track_sequence_and_volume_number(tracks) prep_track_name_length(tracks) context = GrpsIngestionContext(project=project, product=product, tracks=tracks, correlation_id=correlation_id, grps_ingestion_id=grps_ingestion_id, product_type=ProductType[product_type] ) prep_forbidden_sequence(context) return StateMachineSchema().dump(context) def get_project(upc, product, prod_no) -> Project: """Get project data for a given UPC.""" grps_project = execute_snowflake_query(GET_PROJECT_BY_PRODUCT_UPC, {'upc': upc, 'prod_no': prod_no}) if len(grps_project) > 1: raise GetGrpsDataException( f'Expected exactly one project for UPC {upc}, ' f'but got {len(grps_project)}') if not grps_project or len(grps_project) == 0: logger.info(f'Could not get project data for UPC {upc}, ' f'creating using product data') project_artist = Participant( name=product.display_artists[0]) project = Project(name=product.product_name, project_code=product.upc, artist=project_artist) return project grps_project = grps_project[0] rec_project_id = grps_project['REC_PROJECT_ID'] rec_project_participant = execute_snowflake_query( GET_PROJECT_PARTICIPANT_BY_REC_PROJECT_ID, {'rec_project_id': rec_project_id}) if not rec_project_participant or len(grps_project) == 0: raise GetGrpsDataException( f'Éxpected at least one participant ' f'for project {rec_project_id}, but got 0') rec_project_participant = rec_project_participant[0] spotify_uri = rec_project_participant['PARTICIP_SPOTIFY_URI'] if rec_project_participant['PARTICIP_SPOTIFY_URI'] == 'NEW': spotify_uri = None project_artist = Participant( name=rec_project_participant['PARTICIP_FULL_NAME'], spotify_uri=spotify_uri, apple_id=rec_project_participant['APPLE_ARTIST_ID'], ) project = Project(name=grps_project['REC_PROJECT_TITLE'], project_code=grps_project['REC_PROJECT_NUMBER'], artist=project_artist) return project def get_product(upc, prod_no) -> (Product, str): """Get product data for a given UPC.""" grps_product = execute_snowflake_query(GET_PRODUCT_DATA_BY_UPC, {'upc': upc, 'prod_no': prod_no}) if not grps_product or len(grps_product) != 1: raise GetGrpsDataException( f'Expected exactly one product for UPC {upc}, ' f'but got {len(grps_product)}') grps_product = grps_product[0] display_artists = get_product_participants(upc, prod_no) (parent_rep_owner_key, rep_owner_key, repertoire_owner_name) = get_rep_owner_data(upc, prod_no) vendor_id, subaccount_id, do_not_ingest = get_vendor_and_subaccount( parent_rep_owner_key, rep_owner_key) label_name = get_label_name(upc, prod_no) if do_not_ingest and not label_name.startswith( config.PALM_TREE_RECORDS_IMPRINT_PREFIX): raise RepOwnerDoNotIngestException( f'Provided rep owner ({parent_rep_owner_key}, {rep_owner_key}) ' 'is marked as do not ingest in mapping table; ' 'Cannot proceed with ingestion.') catalog_number = grps_product['PROD_NO'] if product_code_is_used_by_different_upc(catalog_number, grps_product['UPC']): catalog_number = catalog_number + '_GRPS' return Product(upc=grps_product['UPC'], product_name=grps_product['PRODUCT_NAME'], grid=grps_product['GRID_NO'], release_date=grps_product['RELEASE_DATE'], release_type=grps_product['RELEASE_TYPE'], catalog_number=catalog_number, sale_start_date=grps_product['SALE_START_DATE'], display_artists=display_artists, repertoire_owner_code=rep_owner_key, repertoire_owner_name=repertoire_owner_name, parent_repertoire_owner_code=parent_rep_owner_key, vendor_id=vendor_id, subaccount_id=subaccount_id, imprint=label_name), grps_product['PRODUCT_TYPE'].upper() def get_product_participants(upc, prod_no): """Get product participants for a given UPC.""" grps_product_track_participants = execute_snowflake_query( GET_PRODUCT_PARTICIPANTS_BY_UPC, {'upc': upc, 'prod_no': prod_no}) participants = [] for part in grps_product_track_participants: spotify_uri = part['PARTICIP_SPOTIFY_URI'] if part['PARTICIP_SPOTIFY_URI'] == 'NEW': spotify_uri = None participant = Participant( name=part['PARTICIP_FULL_NAME'], spotify_uri=spotify_uri, apple_id=part['APPLE_ARTIST_ID'], roles=[part['ROLE']] ) participants.append(participant) return participants def get_tracks(upc, prod_no, release_type) -> list[Track]: """Get tracks for a given UPC.""" grps_product_track_participants = execute_snowflake_query( GET_PRODUCT_TRACKS_BY_UPC, {'upc': upc, 'prod_no': prod_no, 'release_type': release_type}) tracks = [] for grps_track in grps_product_track_participants: participants = get_track_participants(upc, grps_track['ISRC'], prod_no) track = Track( isrc=grps_track['ISRC'], track_name=grps_track['TRACK_NAME'], duration=grps_track['PLAY_TIME'], volume=grps_track['ID_NO'], sequence_number=grps_track['SEQ_NO'], explicit=grps_track['EXPLICIT'], version=grps_track['VERSION'], display_artists=participants ) tracks.append(track) return tracks def get_track_participants(upc, isrc, prod_no) -> list[Participant]: """Get track participants for a given UPC and ISRC.""" grps_track_participants = execute_snowflake_query( GET_PRODUCT_TRACKS_PARTICIPANTS_BY_UPC_ISRC, {'upc': upc, 'isrc': isrc, 'prod_no': prod_no}) participants = [] for grps_participant in grps_track_participants: spotify_uri = grps_participant['PARTICIP_SPOTIFY_URI'] if grps_participant['PARTICIP_SPOTIFY_URI'] == 'NEW': spotify_uri = None participant = Participant( name=grps_participant['PARTICIP_FULL_NAME'], spotify_uri=spotify_uri, apple_id=grps_participant['APPLE_ARTIST_ID'], roles=[grps_participant['ROLE']] ) participants.append(participant) return participants def get_rep_owner_data(upc, prod_no): """Get parent rep owner data.""" grps_parent_rep_owner_key = execute_snowflake_query( GET_PARENT_PARENT_AND_REP_OWNER_KEY_BY_UPC, {'upc': upc, 'prod_no': prod_no}) if len(grps_parent_rep_owner_key) > 1: raise GetGrpsDataException( f'Found ambiguous rep_owner_keys for {upc};') elif not grps_parent_rep_owner_key or len(grps_parent_rep_owner_key) == 0: raise GetGrpsDataException( f'No rep_owner_key found for {upc}') grps_parent_rep_owner_key = grps_parent_rep_owner_key[0] return grps_parent_rep_owner_key['PARENT_REP_OWNER_KEY'], \ grps_parent_rep_owner_key['REP_OWNER_KEY'], grps_parent_rep_owner_key[ 'COMPANY_NAME'] def get_vendor_and_subaccount(parent_rep_owner_key, rep_owner_key): """Get vendor and subaccount for given rep owner keys.""" ddex_mapping_vendor_subaccount = execute_snowflake_query( GET_VENDOR_AND_SUBACCOUNT_BY_PARENT_AND_REP_OWNER_KEY, {'parent_rep_owner_key': parent_rep_owner_key, 'rep_owner_key': rep_owner_key}) if len(ddex_mapping_vendor_subaccount) > 1: raise AmbiguousVendorOrSubaccountMapping( 'Ambiguous vendor/subaccount mapping ' f'parent_rep_owner_key {parent_rep_owner_key} ' f'and rep_owner_key {rep_owner_key}') elif len(ddex_mapping_vendor_subaccount) == 0: logger.warning( f'No vendor/subaccount mapping found for ' f'parent_rep_owner_key {parent_rep_owner_key} ' f'and rep_owner_key {rep_owner_key}; ' 'Defaulting to None for all values.') return None, None, None ddex_mapping_vendor_subaccount = ddex_mapping_vendor_subaccount[0] return ddex_mapping_vendor_subaccount['VENDOR_ID'], \ ddex_mapping_vendor_subaccount['SUBACCOUNT_ID'], \ ddex_mapping_vendor_subaccount['DO_NOT_INGEST'] def get_label_name(upc, prod_no): """Get label related data for a given UPC.""" label = execute_snowflake_query(GET_LABEL_BY_UPC, {'upc': upc, 'prod_no': prod_no}) if not label or len(label) == 0: return None return label[0]['LABEL_NAME'] def product_code_is_used_by_different_upc(product_code, upc): """Check if product_code used by different upc.""" release = execute_snowflake_query( GET_ART_RELATIONS_PRODUCT_BY_PRODUCT_CODE, {'product_code': product_code}) if not release or len(release) == 0: return False return upc != release[0]['UPC']