"""Main Handler.""" import csv from datetime import datetime from io import StringIO import os from types import SimpleNamespace from bulk_metadata_ingester_common.models.base_csv import ( ProjectFieldsModel, TableRowModel) from bulk_metadata_ingester_common.utils.logging import get_current_logger import config from constants.file import ( CSV_DELIMITER, CSV_QUOTECHAR, CSV_VALUE_SEPARATOR, FOLDER_MIME_TYPE, INPUT_FILE_EXT_LIST, # LYRICS_FILE_SUFFIX, # TODO: Add Lyrics MAX_PERFORMER_COUNT, PROCESS_FILE_SUFFIXES) from constants.headers import HEADER_MAPPING from constants.roles import ( BULK_ARTIST_ROLE_TO_ORCHARD_PERFORMER_ROLE as ARTIST_TO_PERFORMER, BULK_BANNED_PERFORMER_ROLES, BULK_RESOURCE_CONTRIBUTOR_ROLE_TO_ORCHARD_PERFORMER_ROLE as CONTRIB_TO_PERFORMER, # noqa FEATURED_PERFORMER_TYPE, PRIMARY_PERFORMER_TYPE ) from constants.territories import EXCLUDE_ONLY from ddex_ingester_common.constants.country_codes import ( ALL_COUNTRY_CODES, WORLDWIDE, ) from ddex_ingester_common.constants.roles import ( DDEX_ARTIST_ROLE_TO_ORCHARD_WRITER as ARTIST_TO_WRITER, # DDEX_RESOURCE_CONTRIBUTOR_ROLE_TO_ORCHARD_PARTICIPANT_ROLE as CONTRIB_TO_PARTICIPANT, # noqa ) from ddex_ingester_common.constants.xmltodict import DDEX_LIST_FIELDS from ddex_ingester_common.lambda_exceptions import CarveoutException from ddex_ingester_common.schemas.ddex_schema import DDEXSchema from lambdacommon.aws import s3 import xmltodict # Get "today" at runtime TODAY = datetime.now().date().isoformat() # Split Constants CONTRIB_TO_PERFORMER_PRIMARY = [ d for d, v in CONTRIB_TO_PERFORMER.items() if v['type'] == 'primary' ] CONTRIB_TO_PERFORMER_NON_FEATURED = [ d for d, v in CONTRIB_TO_PERFORMER.items() if v['type'] != 'primary' ] def handler(event, context): """Convert DDEX to CSV.""" key = event.get('key') bucket = event.get('bucket') correlation_id = event.get('correlation_id') logger = get_current_logger( config.ENVIRONMENT, config.LAMBDA_NAME, logging_level=config.LOGGING_LEVEL, correlation_id=correlation_id) # If single file passed as key (determined by key mimetype) if FOLDER_MIME_TYPE not in s3.head_object(bucket, key)['ContentType']: # Make a key root (no .ext) which is located in the same folder as key. target_key_root = os.path.splitext(key)[0] # Process the file enriched_context, rows = process_single_file(bucket, key, logger) # Write out CSV file write_to_csv(rows, bucket, target_key_root) return enriched_context # If folder passed, get list of all usable files files = [] for k in get_matching_s3_keys(bucket, key, PROCESS_FILE_SUFFIXES): files.append(k) # Parse XML keys to CSV rows, and accumulate them. all_rows = [] for f in files: # Process XML input file if os.path.splitext(f)[1] in INPUT_FILE_EXT_LIST: # TODO: Put context objects into a variable, and add lyrics to the # model if they are available enriched_context, rows = process_single_file(bucket, f, logger) all_rows = all_rows + rows # TODO: Add lyrics full_file_key = 'parsed_csv/ALL_DDEX_TO_CSV_' + TODAY write_to_csv(all_rows, bucket, full_file_key) return enriched_context def check_key_extension(key, logger): """Check file extension is a valid type.""" extension = os.path.splitext(key)[1] if extension not in INPUT_FILE_EXT_LIST: msg = 'S3 metadata file is not a valid file type' logger.warning(msg) raise ValueError(msg) def generate_rows(context, logger): """Generate rows to write to CSV.""" rows = [] # Find all relevant DealTerm dates. release_dates = retrieve_release_dates(context.deals) sales_start_date = release_dates.get('sales_start_date') original_release_date = context.product.original_release_date \ if context.product.original_release_date else sales_start_date earliest_preorder = retrieve_earliest_preorder(context.deals) product_deal_terms = get_product_deal_terms(context) if not product_deal_terms: raise CarveoutException('Product Deal Terms not found.') logger.info(f'Got product deal terms: {product_deal_terms}') carveout_data = format_carveout_data(product_deal_terms, logger) if not carveout_data: raise CarveoutException('Product Deal includes a Takedown.') # get complement territories country_codes = map_country_codes(carveout_data) project_model = ProjectFieldsModel( project_code=context.product.upc, release_name=context.product.product_name, orchard_artist=context.product.display_artist_name, digital_upc=context.product.upc, p_line=context.product.p_line, c_info=context.product.c_line, product_code=context.product.catalog_number, release_meta_language=context.product.metadata_language, artist_url=context.product.artist_profile_page, manufacturers_upc=context.product.proprietary_id, release_version=context.product.product_version, format_type=context.product.release_type, genre=context.product.genres[0].genre, subgenre=context.product.genres[0].subgenre, album_pricing=context.deals[0].deal_terms[0].price_type, # num_agreements=context.product.num_agreements, # These artists requires some sub-processing release_artists_primary_artists=retrieve_display_artists(context.product.display_artists, 'MainArtist'), # noqa release_artists_featurings=retrieve_display_artists(context.product.display_artists, 'FeaturedArtist'), # noqa release_artists_remixers=retrieve_display_artists(context.product.display_artists, 'Remixer'), # noqa release_artists_producers=retrieve_display_artists(context.product.display_artists, 'Producer'), # noqa release_artists_composers=retrieve_display_artists(context.product.display_artists, 'Composer'), # noqa release_artists_orchestras=retrieve_display_artists(context.product.display_artists, 'Orchestra'), # noqa release_artists_ensembles=retrieve_display_artists(context.product.display_artists, 'Ensemble'), # noqa release_artists_conductors=retrieve_display_artists(context.product.display_artists, 'Conductor'), # noqa original_release_date=original_release_date, sale_start_date=sales_start_date, only_include=EXCLUDE_ONLY, territory_iso_codes=country_codes if country_codes else None) for track in context.tracks: # Get all unique roles unique_roles = get_unique_roles(track) # Get track performers performers_list = get_track_performers(track) # Process row my_tuple = TableRowModel( # TODO: We may want to include full path here imprint=context.product.imprint, file_name=track.asset.filename, track_name=track.track_name, track_audio_language=track.lyrics_language, isrc=track.isrc, explicit=track.explicit, volume=track.volume, track_no=track.sequence_number, track_version=track.track_version, track_pricing=get_track_price_tier(track, context.deals), # third_party_publisher=track.third_party_publisher, # only_include=track.only_include, track_p_info=track.p_line, # territory_iso_codes=track.territory_iso_codes, # num_agreements=track.num_agreements, songwriters=get_track_songwriters(unique_roles), track_artist=retrieve_track_artists(track, 'MainArtist'), # noqa: E501 track_artists_featurings=retrieve_track_artists(track, 'FeaturedArtist'), # noqa: E501 track_artists_remixers=retrieve_track_artists(track, 'Remixer'), # noqa: E501 track_artists_producers=retrieve_track_artists(track, 'Producer'), # noqa: E501 track_artists_composers=retrieve_track_artists(track, 'Composer'), # noqa: E501 track_artists_orchestras=retrieve_track_artists(track, 'Orchestra'), # noqa: E501 track_artists_ensembles=retrieve_track_artists(track, 'Ensemble'), # noqa: E501 track_artists_conductors=retrieve_track_artists(track, 'Conductor'), # noqa: E501 publishers=retrieve_track_artists(track, 'MusicPublisher'), itunes_preorder=( 'No Pre-order' if not earliest_preorder else True), itunes_preorder_date=( earliest_preorder.start_date if earliest_preorder else None), preorder_preview=is_previewable(earliest_preorder), # Don't yet know about the data source of these performer_1_type=performers_list[0].type, performer_1_legal_name=performers_list[0].name, performer_1_main_role=performers_list[0].role, performer_2_type=performers_list[1].type, performer_2_legal_name=performers_list[1].name, performer_2_main_role=performers_list[1].role, performer_3_type=performers_list[2].type, performer_3_legal_name=performers_list[2].name, performer_3_main_role=performers_list[2].role, performer_4_type=performers_list[3].type, performer_4_legal_name=performers_list[3].name, performer_4_main_role=performers_list[3].role, performer_5_type=performers_list[4].type, performer_5_legal_name=performers_list[4].name, performer_5_main_role=performers_list[4].role, **project_model._asdict()) rows.append(my_tuple) return rows def write_csv_data(fd, rows): """Write CSV to local file.""" writer = csv.writer( fd, delimiter=CSV_DELIMITER, quotechar=CSV_QUOTECHAR, quoting=csv.QUOTE_ALL ) headers = [ HEADER_MAPPING[d] if d in HEADER_MAPPING else d for d in rows[0]._asdict().keys() ] writer.writerow(headers) for row in rows: writer.writerow(row._asdict().values()) # Return file descriptor, if needed return fd def write_csv_local(rows): """Write data to a local CSV file. Args: rows (list): List of rows to write """ os.makedirs('tmp/', exist_ok=True) with open('tmp/test_output.csv', 'w', newline='') as csvfile: write_csv_data(csvfile, rows) def write_csv_s3(rows, bucket, key): """Write data to CSV on S3. Args: rows (list): List of rows to write. """ csv_buffer = write_csv_data(StringIO(), rows) target_key = key + '.csv' s3.resource.Object(bucket, target_key).put(Body=csv_buffer.getvalue()) def write_to_csv(rows, bucket, key): """Write data to CSV. Args: rows (list): List of rows to write """ write_csv_local(rows) write_csv_s3(rows, bucket, key) # TODO: Add missing data to context def enrich_context(context): """Enrich the context with Orchard-specific information.""" # Do nothing at the moment # These will likely be hydrated ? # publishers=track.publishers, # orch_label_id=track.orch_label_id, # subaccount_name=context.product.subaccount_id, # track_lyrics=track.lyrics, # ownership_for_this_sound_recording=track.ownership_rights, # noqa # country_of_recording=track.recording_country_code, # nationality_of_original_copyright_owner=track.copyright_owner_country, # noqa return context def retrieve_track_artists(track, role): """Retrieve track artist and contributor names and concatenate together.""" names = retrieve_display_artists(track.display_artists, role) if not names: names = retrieve_display_artists(track.resource_contributors, role) return names def retrieve_display_artists(display_artists, role): """Retrieve display artist names and concatenate together.""" names = '' for artist in display_artists: if role in artist.roles: if artist.name: if not names: names = artist.name else: names += f'{CSV_VALUE_SEPARATOR}{artist.name}' return names # The next 3 functions were copied from: # lambda-ddex-ingester.lambda.set_preorder_info.index def retrieve_earliest_preorder(deals): """Retrieve earliest pre-order deal terms object.""" earliest_deal_terms = None earliest_date_seen = None for deal in deals: for term in deal.deal_terms: if not term.pre_order: continue start_date = term.start_date start_date_time = term.start_date_time early_date = None early_date = get_earliest_date(start_date, start_date_time) if not earliest_date_seen and early_date: earliest_date_seen = early_date earliest_deal_terms = term if early_date and early_date < earliest_date_seen: earliest_date_seen = early_date earliest_deal_terms = term return earliest_deal_terms def is_previewable(earliest_preorder) -> bool: """Set preview Y/N for product.""" if not earliest_preorder: return False pre_order_date = get_earliest_date( earliest_preorder.start_date, earliest_preorder.start_date_time) if not earliest_preorder.clip_preview_date: return False clip_preview_date = earliest_preorder.clip_preview_date return clip_preview_date <= pre_order_date def get_earliest_date(one: str, two: str) -> str: """Not-safe date comparison between two strings.""" if not one or not two: return one or two if one > two: return two return one def process_single_file(bucket, key, logger): """Process a single XML file into a context object and a list of rows. Args: bucket (str): The S3 bucket containing the key. key (str): The S3 key to process. Returns: tuple (dict, list): The context (dict), and the flattened rows from the file. """ # Peel off file location info # Validate input file extenison # check_key_extension(key) # Grab S3 XML object stream fd = s3.get_object(bucket, key)['Body'] # Convert S3 XML stream to dict doc = xmltodict.parse(fd.read(), force_list=DDEX_LIST_FIELDS) try: context = DDEXSchema().load(doc) except (TypeError, ValueError) as err: file = os.path.split(key)[1] logger.warning(f'{file} is missing a component. Cannot perform ' f'DDEXSchema.load(). \n{str(err)}') return {}, [] # Find what language this content is define_content_languages(context, doc) try: check_territory_codes(context, logger) except ValueError as err: file = os.path.split(key)[1] logger.warning(f'{file} is NOT WORLDWIDE. Cannot perform ' f'DDEXSchema.load(). \n{str(err)}') # return {}, [] # TODO: Add missing data enriched_context = enrich_context(context) # Genrate CSV rows try: rows = generate_rows(enriched_context, logger) except (TypeError, AttributeError, ValueError) as err: file = os.path.split(key)[1] logger.warning( f'{file} is missing a component. Cannot convert to rows. ' f'(generate_rows()) \n{str(err)}') return {}, [] except (CarveoutException) as err: file = os.path.split(key)[1] logger.warning( f'{file} has a carveout issue. Cannot convert to rows. ' f'(generate_rows()) \n{str(err)}') return {}, [] return enriched_context, rows # TODO: Inject lyrics into CSV def add_lyrics(bucket, key): """Process a lyrics .TXT file, and return a properly escaped str.""" pass def get_matching_s3_keys(bucket, prefix='', suffix=''): """Yield the keys in an S3 bucket as a generator. From https://alexwlchan.net/2017/listing-s3-keys/ Args: bucket (str): Name of the S3 bucket. prefix (str, optional): Only fetch keys that start with this prefix. Defaults to ''. suffix (str, optional): Only fetch keys that end with this suffix. Defaults to ''. Yields: str: a key from the bucket that meets the prefix / suffix criteria. """ kwargs = {'Bucket': bucket} # If the prefix is a single string (not a tuple of strings), we can # do the filtering directly in the S3 API. if isinstance(prefix, str): kwargs['Prefix'] = prefix while True: # The S3 API response is a large blob of metadata. # 'Contents' contains information about the listed objects. resp = s3.client.list_objects_v2(**kwargs) for obj in resp['Contents']: key = obj['Key'] if key.startswith(prefix) and key.endswith(suffix): yield key # The S3 API is paginated, returning up to 1000 keys at a time. # Pass the continuation token into the next response, until we # reach the final page (when this field is missing). try: kwargs['ContinuationToken'] = resp['NextContinuationToken'] except KeyError: break def get_track_price_tier(track, deals): """Get a track's price type from the deal.""" for deal in deals: for deal_term in deal.deal_terms: if deal_term.price_type and \ track.release_reference in deal.release_references: return deal_term.price_type return None def define_content_languages(context, doc): """Define what language the Release is.""" fallback = doc['ern:NewReleaseMessage']['@LanguageAndScriptCode'] or None if not context.product.metadata_language: # If no metadata language, default to main level language. context.product.metadata_language = fallback def get_unique_roles(track): """Get a dict with every role along with the participant's name. Modified from: https://github.com/theorchard/lambda-ddex-ingester/blob/829ede51e733b3bd4a316e2b6a5f8a15fa568957/lambda/set_track_metadata/index.py#L176 # noqa """ unique_roles = {} for artist in track.display_artists: for role in artist.roles: if unique_roles.get(role): unique_roles[role] = \ CSV_VALUE_SEPARATOR.join([unique_roles[role], artist.name]) else: unique_roles[role] = artist.name for contributor in track.resource_contributors: for role in contributor.roles: if unique_roles.get(role): unique_roles[role] = \ CSV_VALUE_SEPARATOR.join([unique_roles[role], contributor.name]) # noqa else: unique_roles[role] = contributor.name return unique_roles def get_track_songwriters(unique_roles, allow_multiple=True): """Get a track's songwriters. Args: unique_roles (dict): A mapping of role names to artist names allow_multiple (bool, optional): Allow concatenation of multiple values. Defaults to True. Returns: str: The songwriter(s) as a string. """ songwriters = [] for role, name in unique_roles.items(): if role in ARTIST_TO_WRITER.keys(): if name not in songwriters: songwriters.append(str(name)) # Cast to coerce weird names if not len(songwriters): return '' elif allow_multiple: return CSV_VALUE_SEPARATOR.join(songwriters) else: return songwriters[0] def get_track_performers(track): """Get track performers for track.""" performers = { 'primary': [], 'featured': [] } # Start with contributors for contrib in track.resource_contributors: for role in contrib.roles: # Check for PRIMARY performer role if role in CONTRIB_TO_PERFORMER_PRIMARY: performers['primary'].append( SimpleNamespace( name=contrib.name, role=role, type=PRIMARY_PERFORMER_TYPE ) ) # Check for FEATURED role (Is it called 'non-featured' in DDEX...?) elif role in CONTRIB_TO_PERFORMER_NON_FEATURED: performers['featured'].append( SimpleNamespace( name=contrib.name, role=role, type=FEATURED_PERFORMER_TYPE ) ) # elif role in CONTRIB_TO_PARTICIPANT: # performers['featured'].append( # SimpleNamespace( # name=contrib.name, # role=role, # type=FEATURED_PERFORMER_TYPE # ) # ) else: # DEBUG if role not in BULK_BANNED_PERFORMER_ROLES: performers['featured'].append( SimpleNamespace( name=contrib.name, role=role, type='FAILED_TO_MAP_CONTRIB' ) ) # Grab all the names to eliminate dupes. primary_names = [a.name for a in performers['primary']] featured_names = [a.name for a in performers['featured']] # Backup check through display_artists for artist in track.display_artists: for role in artist.roles: # Check for perfomer role (only primary exist in these) if role in ARTIST_TO_PERFORMER.keys(): if artist.name not in primary_names: # De-Dupe performers['primary'].append( SimpleNamespace( name=artist.name, role=role, type=PRIMARY_PERFORMER_TYPE ) ) else: # DEBUG if artist.name not in featured_names: # De-Dupe if role not in BULK_BANNED_PERFORMER_ROLES: performers['featured'].append( SimpleNamespace( name=artist.name, role=role, type='FAILED_TO_MAP_ARTIST' ) ) # Reverse order to give contribs priority over display_artists performers['primary'].reverse() performers['featured'].reverse() # Make a nice neat list for processing performer_list = [] for _ in range(0, MAX_PERFORMER_COUNT): # Max 5 right now. Priority order if len(performers['primary']): performer_list.append(performers['primary'].pop()) elif len(performers['featured']): performer_list.append(performers['featured'].pop()) else: # Flood with blank object. performer_list.append( SimpleNamespace(name='', role='', type='')) return performer_list def check_territory_codes(context, logger): """Check if we have any non-Worldwide territory codes.""" for deal in context.deals: for term in deal.deal_terms: territory_set = set(term.territories) if territory_set != {'Worldwide'}: logger.warning(f'Expected only Worldwide territories, found {territory_set}.') # noqa raise ValueError('Carveouts found for product.') def retrieve_release_dates(deals): """Retrieve the latest start date from deal terms.""" # Only keep R0 or R1 release deal, in that order. filtered_deal = filter_release_deals(deals) # We have a filtered deal, weed out pre_orders filtered_deal = filter_preorder_deals(filtered_deal) # Find earliest sales start date and set timed release date sales_start_date, timed_release_date = retrieve_deal_dates( filtered_deal) # Gather territory dates that are not the same as sales start date territory_dates = retrieve_territory_with_differing_dates( filtered_deal, sales_start_date) return { 'sales_start_date': sales_start_date, 'timed_release_date': timed_release_date, 'territory_dates': territory_dates, } def filter_release_deals(deals): """Try to filter only for R0 then R1 if available. Otherwise fail.""" filtered_deals = [deal for deal in deals if 'R0' in deal.release_references] if not filtered_deals: filtered_deals = [deal for deal in deals if 'R1' in deal.release_references] if not filtered_deals: raise ValueError('No valid release deals found.') return filtered_deals[0] def filter_preorder_deals(deal): """Filter out preorder deal terms from release deal.""" # We have a filtered deal, weed out pre_orders deal.deal_terms = [deal_term for deal_term in deal.deal_terms if not deal_term.pre_order] if not deal.deal_terms: raise ValueError('Only pre-order deals exist in R0 or R1.') return deal def retrieve_deal_dates(deal): """Retrieve earliest sales_start_date and set timed_release_date.""" sales_start_date = None timed_release_date = None for deal_term in deal.deal_terms: start_date = deal_term.start_date start_date_time = deal_term.start_date_time if not start_date and not start_date_time: continue if not start_date: # Strip time from datetime start_date = start_date_time[0:10] # A date time with format '00:00:00.000Z' would indicate a timed # release and '00:00:00.000' would indicate a non-timed release. # Timed release date is the latest deal term start date set_timed_release_date = ( start_date_time.endswith('Z') and ( not timed_release_date or timed_release_date < start_date_time) # noqa ) if set_timed_release_date: timed_release_date = start_date_time if not sales_start_date or sales_start_date > start_date: sales_start_date = start_date return sales_start_date, timed_release_date def retrieve_territory_with_differing_dates( filtered_deal, sales_start_date): """Retrieve territories when start date != sales start date.""" territory_dates = {} for deal_term in filtered_deal.deal_terms: start_date = deal_term.start_date start_date_time = deal_term.start_date_time if not start_date and not start_date_time: continue if not start_date: # Strip time from datetime start_date = start_date_time[0:10] if start_date == sales_start_date: continue if deal_term.territories: for territory in deal_term.territories: territory_date = territory_dates.get(territory) if not territory_date or territory_date < start_date: territory_dates[territory] = start_date return territory_dates def get_product_deal_terms(context): """Get deal term for R0 if not then R1 else None.""" for deal in context.deals: if 'R0' in deal.release_references: return deal.deal_terms if 'R1' in deal.release_references: return deal.deal_terms return None def format_carveout_data(deal_terms, logger): """Extract territories from product deal terms.""" territories = [] for deal_term in deal_terms: if deal_term.takedown: return None # If we encounter a deal EndDate in the past, skip this term. end_date = deal_term.end_date_time or deal_term.end_date if end_date: formatted_date = datetime.strptime(end_date[:10], '%Y-%M-%d') if formatted_date < datetime.now(): end_date_territories = deal_term.territories logger.info(f'End date found for deal term: {end_date} ' f'Skipping these territories: {end_date_territories}') # noqa continue if deal_term.territories: if WORLDWIDE in deal_term.territories: territories = ALL_COUNTRY_CODES else: territories.extend(deal_term.territories) elif deal_term.excluded_territories: territories.extend( list(set( ALL_COUNTRY_CODES) - set(deal_term.excluded_territories))) return territories def map_country_codes(country_codes): """Map country codes to Orchard.""" mapped_country_codes = set() for country_code in country_codes: mapped_country_codes.add(country_code) return sorted(list(set(ALL_COUNTRY_CODES) - mapped_country_codes))