"""Lambda create_hfa_request_files function module.""" import os import csv import json import random import string import datetime from typing import Any, Dict, List, Optional, Tuple import tempfile from sentry_sdk import capture_exception import pandas as pd import config from lambdacommon.common_config import logger from lambdacommon.aws import s3 from lambdacommon.util import init_sentry_for_lambda from src.constants import request_file as file from src.constants import track_fields as field from src.constants import common init_sentry_for_lambda() pd.options.display.max_colwidth = 250 def is_null(value: Optional[str]) -> bool: """ Check if the value is NaN in a pandas DataFrame cell. Args: value (Optional[str]): The value from a pandas DataFrame cell to be checked, which could be a string or NaN. Returns: bool: True if the value is NaN, False otherwise. """ return pd.isna(value) def get_length_string(value: str, length: int) -> str: """ Trims the input string to the specified length if the string's length exceeds the specified value. Args: value (str): The string to be trimmed. length (int): The maximum allowable length for the string. Returns: str: The original string if its length is less than or equal to the specified length, otherwise the trimmed string. """ return value[:length] if len(value) > length else value def safe_str(value: Optional[Any]) -> str: """Safely converts the input value to a string. If the value is null (None or NaN), it returns an empty string. Also removes tab characters and strips leading/trailing whitespaces. Args: value (Optional[Any]): The value to be converted to a string. Returns: str: The safely converted string, or an empty string if the value is null. """ return '' if is_null(value) else str(value).replace('\t', '').strip() def safe_trim(value: Optional[Any], length: int) -> str: """Safely trims the input value to a specified length. If the value is null (None or NaN), it returns an empty string. The value is first converted to a string before trimming. Args: value (Optional[Any]): The value to be trimmed. length (int): The maximum length to which the string should be trimmed. Returns: str: The safely trimmed string, or an empty string if the value is null. """ return '' if is_null(value) else get_length_string(str(value), length) def transform( row12: pd.DataFrame, track_info: pd.DataFrame, resubmit_licenses: pd.DataFrame, current_date: str ) -> pd.DataFrame: """ Transform multiple input DataFrames into a single formatted DataFrame. Args: row12 (pd.DataFrame): HFA orchard track licenses. track_info (pd.DataFrame): Track metadata. resubmit_licenses (pd.DataFrame): Filtered licenses needing resubmission. current_date (str): Current date in YYYYMMDD format. Returns: pd.DataFrame: Transformed DataFrame ready for output or further processing. """ logger.info('Started transforming track and license data') data = [] try: for _, row in track_info.iterrows(): select_row12 = row12.loc[row12[field.ORCHARD_TRACK_ID] == int(row[field.TRACK_ID])] for _, select_row12_row in select_row12.iterrows(): temp = dict.fromkeys(file.FIELDS, '') select_resubmit_licenses = resubmit_licenses.loc[ resubmit_licenses[field.ID] == int(select_row12_row[field.ID]) ] config_code = str(select_row12_row[field.HFA_CONFIGURATION_CODE]).upper() is_rt = config_code == file.RT is_sp = config_code == file.SP var_agreement_code = file.RGT if is_rt else file.SSA temp[file.HFA_AGREEMENT_CODE] = var_agreement_code temp[file.MANUFACTURER_NUMBER] = file.ORCHARD_MANUFACTURER_NUMBER temp[file.TRANSACTION_DATE] = int(current_date) temp[file.MANUFACTURER_REQUEST_NUMBER] = int(select_row12_row[field.ID]) temp[file.LABEL_NAME] = safe_str(row[field.LABEL_NAME]) temp[file.ISRC_CODE] = str(row[field.ISRC]) temp[file.PLAYING_TIME_MINUTES] = int(row[field.LENGTH_MINUTE]) temp[file.PLAYING_TIME_SECONDS] = int(row[field.LENGTH_SECONDS]) temp[file.ARTIST] = safe_trim(row[field.ARTIST_NAME], 200) temp[file.SONG_TITLE] = safe_trim(row[field.TRACK_NAME], 200) temp[file.AKA_SONG_TITLE] = safe_trim(row[field.TRACK_NAME], 200) temp[file.COMPOSERS] = safe_trim(row[field.WRITERS], 200) temp[file.PUBLISHER_NAME] = safe_trim(row[field.PUBLISHER_NAME], 60) temp[file.CATALOG_NUMBER] = safe_trim(row[field.VENDOR_CATALOG_NUMBER], 15) temp[file.ALBUM_TITLE] = safe_trim(row[field.RELEASE_NAME], 200) temp[file.UPC_CODE] = str(row[field.DISPLAY_UPC])[0:16] temp[file.CONFIGURATION_CODES] = select_row12_row[field.HFA_CONFIGURATION_CODE] temp[file.LICENSE_TYPE] = 'G' if is_rt or is_sp else 'D' temp[file.SERVER_FIXATION_DATE] = int(row[field.INGESTION_DATE].replace('-', '')) temp[file.RATE_CODE] = 'Y' if var_agreement_code.upper() == file.RGT else 'S' field_id = select_row12_row[field.ID] resubmit_field_id = select_resubmit_licenses[field.ID].to_string(index=False) hfa_process_code_1 = 'R' if field_id == resubmit_field_id else '' temp[file.HFA_REQUEST_PROCESS_CODE_1] = '' if select_resubmit_licenses.empty else hfa_process_code_1 temp[file.USER_DEFINED_1] = 'N2' temp[file.USER_DEFINED_2_TRACK_ID] = str(select_row12_row[field.ORCHARD_TRACK_ID]) temp[file.USER_DEFINED_3_COPYRIGHT_OFFICE_REQUEST] = 'FALSE' temp[file.USER_DEFINED_4_DISTRIBUTION_DATE] = str(row[field.DISTRIBUTION_DATE].replace('-', '')) temp[file.USER_DEFINED_5_TRACK_PUBLISHER_ID] = ( '' if is_null(row[field.TRACK_PUBLISHER_ID]) else str(row[field.TRACK_PUBLISHER_ID]).split('.')[0] ) temp[file.USER_DEFINED_8_PRIORITY_CODE] = safe_trim(row[field.PRIORITY], 18) data.append(temp) if not data: logger.warning('No data was transformed. The resulting DataFrame will be empty.') return pd.DataFrame(columns=file.FIELDS) list_df = pd.DataFrame(data, columns=file.FIELDS) list_df = list_df.apply(lambda x: x.str.replace('\t', '') if x.dtype == 'object' else x) logger.info(f'Successfully transformed data: {len(list_df)} rows generated') return list_df except Exception as e: logger.exception(f'Failed during data transformation: {e}') raise def generate_flat_track_data(track: Dict[str, Any], publisher: Optional[Dict[str, Any]] = None) -> Dict[str, Any]: """ Convert nested track and publisher dictionaries into a flat dictionary. Args: track (Dict[str, Any]): Track details dictionary. publisher (Dict[str, Any], optional): Publisher details dictionary. Returns: Dict[str, Any]: Flattened dictionary of track and publisher details. """ try: publisher = publisher or {} temp: Dict[str, Any] = {} temp[field.TRACK_ID] = track[field.TUID] temp[field.LABEL_NAME] = track[field.PRODUCT][field.LABEL][field.NAME] temp[field.ISRC] = track[field.ISRC] temp[field.LENGTH_MINUTE] = track[field.DURATION_MINUTES] temp[field.LENGTH_SECONDS] = track[field.DURATION_SECONDS] temp[field.ARTIST_NAME] = track[field.PRODUCT][field.PROJECT][field.LABEL_PARTICIPANT][field.NAME] temp[field.TRACK_NAME] = track[field.GRAPHQL_TRACK_NAME] temp[field.PUBLISHER_NAME] = publisher[field.NAME] if bool(publisher) else None temp[field.VENDOR_CATALOG_NUMBER] = track[field.PRODUCT][field.GRAPHQL_VENDOR_CATALOG_NUMBER] temp[field.RELEASE_NAME] = track[field.PRODUCT][field.PRODUCT_NAME] temp[field.DISPLAY_UPC] = track[field.PRODUCT][field.GRAPHQL_DISPLAY_UPC] temp[field.INGESTION_DATE] = track[field.PRODUCT][field.INGESTION_COMPLETED] temp[field.DISTRIBUTION_DATE] = track[field.PRODUCT][field.GRAPHQL_DISTRIBUTION_DATE] temp[field.TRACK_PUBLISHER_ID] = publisher[field.ID] if bool(publisher) else None temp[field.PRIORITY] = track[field.PRODUCT][field.PRIORITY] temp[field.WRITERS] = ','.join( participant[field.PARTICIPANT][field.NAME] for participant in track[field.PARTICIPATIONS] ) return temp except Exception: logger.exception(f"Failed to create dictionary for track ID: {track.get(field.TUID, 'unknown')}") raise def _is_valid_track_data(data: Dict[str, Any]) -> bool: """Check if track data contains all required fields with no None values. Args: data (Dict[str, Any]): Track data. Returns: bool: True if data is valid, False otherwise. """ return None not in ( data.get(field.LENGTH_MINUTE), data.get(field.LENGTH_SECONDS), data.get(field.INGESTION_DATE) ) def track_info_json_to_df(tracks: List[Dict]) -> pd.DataFrame: """Convert list of track details json into dataframe. Args: tracks (list): list of track's info. Returns: dataframe: Track info dataframe. """ try: logger.info('Converting track info JSON into DataFrame') track_data = [] for track in tracks: publishers = track[field.PUBLISHING][field.PUBLISHERS] if publishers: for publisher in publishers: data = generate_flat_track_data(track, publisher) if _is_valid_track_data(data): track_data.append(data) else: data = generate_flat_track_data(track) if _is_valid_track_data(data): track_data.append(data) if not track_data: logger.warning('No valid track data found. Returning empty DataFrame.') return pd.DataFrame(columns=field.TRACK_FIELDS) track_info_df = pd.DataFrame(track_data, columns=field.TRACK_FIELDS) track_info_df.sort_values(by=[field.TRACK_ID, field.TRACK_PUBLISHER_ID], inplace=True) logger.info(f'Successfully transformed track info JSON to DataFrame with {len(track_info_df)} rows') return track_info_df except Exception as e: logger.exception(f'Failed to convert track info JSON into DataFrame: {e}') raise def generate_random_string(length: int = 5) -> str: """ Generate a random alphanumeric string of the given length. Args: length (int): The length of the random string to generate (default is 5). Returns: str: A random alphanumeric string of the specified length. """ return ''.join( random.SystemRandom().choice(string.ascii_letters + string.digits) for _ in range(length) ) def generate_file_name(manufacturer_number: str, file_type: str, current_date: str, random_string: str) -> str: """ Generate a formatted filename for the HFA request. Args: manufacturer_number (str): The manufacturer number to include in the filename. file_type (str): The file type (e.g., "licenses" or "metadata"). current_date (str): The current date to include in the filename, usually in YYYYMMDD format. random_string (str): A random string to be appended to the filename for uniqueness. Returns: str: The formatted filename as a string. """ return common.FILE_NAME_TEMPLATE.format( manufacturer_number=manufacturer_number, file_type=file_type, current_date=current_date, random_string=random_string ) def upload_to_s3(file_name: str, data_df: pd.DataFrame) -> str: """Convert DataFrame to CSV and upload to S3. Args: file_name (str): Name of the file to upload. data_df (pd.DataFrame): Data to write to the file. Returns: str: Uploaded file name. """ s3_file_path = f'{config.S3_REQUEST_FILE_DIR}{file_name}' full_s3_path = f's3://{config.S3_BUCKET_NAME}/{s3_file_path}' logger.info('Preparing to upload HFA request file to S3: %s', file_name) try: with tempfile.TemporaryDirectory() as tmpdir: local_file_path = os.path.join(tmpdir, file_name) data_df.to_csv( local_file_path, index=False, header=False, sep='\t', na_rep='', quoting=csv.QUOTE_NONE, escapechar='\\' ) logger.info('Uploading file from S3: %s to %s', local_file_path, full_s3_path) s3.upload_file_to_s3( config.S3_BUCKET_NAME, local_file_path, s3_file_path ) logger.info('File successfully uploaded to S3: %s', full_s3_path) return s3_file_path except Exception as e: logger.exception('Failed to upload file to S3. Error: %s', str(e)) raise def generate_and_upload_file(agreement_code: str, filtered_df: pd.DataFrame, current_date: str) -> dict: """Generate and upload a file for a given agreement code and DataFrame. Args: agreement_code (str): Agreement code to filter the DataFrame. filtered_df (pd.DataFrame): The DataFrame containing the data to filter and upload. current_date (str): Current date for file naming. Returns: dict: Contains 'file_name' and 's3_file_path' for further use. """ try: random_string = generate_random_string() file_name = generate_file_name(file.ORCHARD_MANUFACTURER_NUMBER, agreement_code, current_date, random_string) s3_file_path = upload_to_s3(file_name, filtered_df) return { 'file_name': file_name, 's3_file_path': s3_file_path } except Exception as e: logger.exception(f'Failed to generate and upload {agreement_code} file: {e}') raise def load_generate_hfa_tmp_files( object_key_map: Dict[str, str] ) -> Tuple[Optional[pd.DataFrame], Optional[pd.DataFrame]]: """ Download and parse temporary HFA files from S3. Args: object_key_map (Dict[str, str]): Dict of file type to S3 key (e.g., {"hfa_orchard_track_licenses_file": "file.json"}). Returns: Tuple[Optional[pd.DataFrame], Optional[pd.DataFrame]]: - track_info_df: DataFrame parsed from JSON. - hfa_orchard_track_licenses_df: DataFrame parsed from JSON. """ try: track_info_df = None hfa_orchard_track_licenses_df = None with tempfile.TemporaryDirectory() as tmpdir: for key_name, key in object_key_map.items(): s3_key = config.S3_TMP_FILE_DIR + key local_path = os.path.join(tmpdir, key) if not s3.object_exists(config.S3_BUCKET_NAME, s3_key): logger.error(f'S3 object not found: s3://{config.S3_BUCKET_NAME}/{s3_key}') raise Exception(f'S3 object not found: {s3_key}') with open(local_path, 'wb') as f: s3.download_file_object(config.S3_BUCKET_NAME, s3_key, f) logger.info(f'Download complete: {local_path}') if key_name == common.HFA_ORCHARD_TRACK_LICENSES_FILE: hfa_orchard_track_licenses_df = pd.read_json(local_path) logger.info(f'Parsed {len(hfa_orchard_track_licenses_df)} records from HFA Orchard licenses file') elif key_name == common.PENDING_REQUEST_TRACK_IDS_DATA_FILE: with open(local_path, encoding='utf8') as f: data = json.load(f) track_info_df = track_info_json_to_df(data) logger.info(f'Parsed {len(track_info_df)} records from track info file') logger.info('Successfully loaded generated HFA tmp files') return hfa_orchard_track_licenses_df, track_info_df except Exception as e: logger.exception(f'Failed to load and generate HFA tmp files: {e}') raise def handler(event, context): """Lambda entry point.""" try: logger.info('Lambda execution started') result = { common.STATUS: event.get(common.STATUS, common.OK), common.GENERATE_HFA_TMP_FILES: event.get(common.GENERATE_HFA_TMP_FILES, {}), common.GENERATE_HFA_TMP_FILES_LENGTH: event.get(common.GENERATE_HFA_TMP_FILES_LENGTH, 0), common.CREATE_HFA_REQUEST_FILES: event.get(common.CREATE_HFA_REQUEST_FILES, {}), common.CREATE_HFA_REQUEST_FILES_LENGTH: event.get(common.CREATE_HFA_REQUEST_FILES_LENGTH, 0), common.FILES_TO_CLEANUP: event.get(common.FILES_TO_CLEANUP, []) } if not result[common.GENERATE_HFA_TMP_FILES]: logger.info('No temporary files provided. Skipping file generation.') return result current_date = datetime.datetime.today().strftime('%Y%m%d') request_files: Dict[str, str] = {} files_to_cleanup: List = result[common.FILES_TO_CLEANUP] hfa_orchard_track_licenses_df, track_info_df = load_generate_hfa_tmp_files( result[common.GENERATE_HFA_TMP_FILES] ) resubmit_orchard_track_licenses_df = hfa_orchard_track_licenses_df.query('state == "resubmit"') r_df = transform( hfa_orchard_track_licenses_df, track_info_df, resubmit_orchard_track_licenses_df, current_date ) ssa_df = r_df.query('hfa_agreement_code == "SSA"') if not ssa_df.empty: logger.info(f'Generating SSA file with {len(ssa_df)} rows') ssa_file_info = generate_and_upload_file(file.SSA, ssa_df, current_date) request_files[common.SSA_REQUEST_FILE] = ssa_file_info['file_name'] files_to_cleanup.append(ssa_file_info['s3_file_path']) rgt_df = r_df.query('hfa_agreement_code == "RGT"') if not rgt_df.empty: logger.info(f'Generating RGT file with {len(rgt_df)} rows') rgt_file_info = generate_and_upload_file(file.RGT, rgt_df, current_date) request_files[common.RGT_REQUEST_FILE] = rgt_file_info['file_name'] files_to_cleanup.append(rgt_file_info['s3_file_path']) result[common.CREATE_HFA_REQUEST_FILES] = request_files result[common.CREATE_HFA_REQUEST_FILES_LENGTH] = len(request_files) result[common.FILES_TO_CLEANUP] = files_to_cleanup return result except Exception as e: capture_exception(e) logger.exception(str(e)) raise