"""Logic for handling HFA-related data processing and filtering.""" from typing import Dict, List import pandas as pd from src.utils import dataframes from src.models import ows_carveouts_python, ows_contracts, ows_product, ows_royalties from src.utils import constants from lambdacommon.common_config import logger import asyncio def fetch_and_load_abacus_active_contracts() -> pd.DataFrame: """Fetch JSON from API and load it into a DataFrame. Returns: pd.DataFrame: DataFrame containing mechanical deduction contracts. """ try: ows_royalties_data = ows_royalties.get_all_active_contract_mechanical_deductions() contract_df = pd.DataFrame(ows_royalties_data) contract_df = contract_df.astype( {'account_id': 'int64', 'contract_id': 'int64', 'term_type': 'string'} ) return contract_df except Exception as e: raise Exception( f'Unexpected error while fetching Abacus Active Contracts {e}' ) from e def parse_mechanical_type_to_list(mechanical_type): """ Convert a mechanical type mechanical_type into a consistent list format. Handles NaN by returning an empty list. If `mechanical_type` is a string, it's split by commas and stripped of whitespace. If `mechanical_type` is already a list, it's returned as is. Unexpected types also result in an empty list. Args: mechanical_type (str, list, float): The mechanical type input. Returns: list: A list of mechanical types, or an empty list if not applicable. """ if isinstance(mechanical_type, str): return [v.strip() for v in mechanical_type.split(',')] if isinstance(mechanical_type, list): return mechanical_type return [] def filter_eligible_hfa_digital_physical_tracks( raw_digital_physical_df: pd.DataFrame, oa_active_contracts_df: pd.DataFrame, abacus_active_contracts_df: pd.DataFrame, carveouts_df: pd.DataFrame ) -> pd.DataFrame: """ Filter digital and physical release tracks eligible for HFA licensing. Tracks must: - Have either an active OA contract (with valid publishing types) OR an active Abacus contract (with valid mechanical types for the transaction type). - Exclude U.S. territory carve-outs (country ID '1'). - Match HFA configuration codes based on transaction type. Args: raw_digital_physical_df (pd.DataFrame): Raw digital/physical track data. oa_active_contracts_df (pd.DataFrame): OA active contracts. abacus_active_contracts_df (pd.DataFrame): Abacus active contracts. carveouts_df (pd.DataFrame): US carveouts data. Returns: pd.DataFrame: Filtered DataFrame with ['track_id', 'hfa_configuration_code']. """ if raw_digital_physical_df.empty: return pd.DataFrame(columns=['track_id', 'hfa_configuration_code']) merged_tracks_df = raw_digital_physical_df.merge( oa_active_contracts_df, how='left', on='vendor_id' ) merged_tracks_df = merged_tracks_df.merge( carveouts_df, how='left', on='release_id' ) merged_tracks_df = merged_tracks_df.merge( abacus_active_contracts_df, how='left', on=['track_id', 'track_isrc', 'upc', 'vendor_id'], suffixes=('', '_abacus'), ) merged_tracks_df['mechanical_type'] = merged_tracks_df['mechanical_type'].apply( parse_mechanical_type_to_list ) merged_tracks_df['has_valid_oa_contract'] = ( (merged_tracks_df['trans_type'] == 'digital') & (merged_tracks_df['digital_track']) ) | ( (merged_tracks_df['trans_type'] == 'physical') & (merged_tracks_df['physical_track']) ) merged_tracks_df['has_valid_abacus_contract'] = ( (merged_tracks_df['trans_type'] == 'digital') & ( merged_tracks_df['mechanical_type'].map(lambda m: 'digital' in m)) ) | ( (merged_tracks_df['trans_type'] == 'physical') & ( merged_tracks_df['mechanical_type'].map(lambda m: 'physical' in m)) ) merged_tracks_df = merged_tracks_df[ (merged_tracks_df['has_valid_oa_contract'] | merged_tracks_df['has_valid_abacus_contract']) & ~( merged_tracks_df['has_us_carveout']) ] if merged_tracks_df.empty: return pd.DataFrame(columns=['track_id', 'hfa_configuration_code']) merged_tracks_df = merged_tracks_df.reset_index().rename(columns={'index': '_original_order'}) filtered_tracks_df = merged_tracks_df.sort_values(by='_original_order').drop( columns=['_original_order'] ) filtered_tracks_df = filtered_tracks_df[filtered_tracks_df['track_id'].notnull()] filtered_tracks_df = filtered_tracks_df[['track_id', 'hfa_configuration_code']] return filtered_tracks_df def filter_eligible_hfa_ringtones( raw_ringtones_df: pd.DataFrame, oa_active_contracts_df: pd.DataFrame, abacus_active_contracts_df: pd.DataFrame, carveouts_df: pd.DataFrame ) -> pd.DataFrame: """Filter ringtone tracks eligible for HFA licensing. Tracks must: - Have either an active OA contract (with valid ringtone publishing types) OR an active Abacus contract (with valid mechanical types). - Exclude U.S. territory carve-outs (country ID '1'). Args: raw_ringtones_df (pd.DataFrame): Raw ringtone track data. oa_active_contracts_df (pd.DataFrame): OA active contracts. abacus_active_contracts_df (pd.DataFrame): Abacus active contracts. carveouts_df (pd.DataFrame): US carveouts data. Returns: pd.DataFrame: Filtered DataFrame with ['track_id', 'hfa_configuration_code']. """ if raw_ringtones_df.empty: return pd.DataFrame(columns=['track_id', 'hfa_configuration_code']) merged_ringtones_df = raw_ringtones_df.merge(oa_active_contracts_df, how='left', on='vendor_id') merged_ringtones_df = merged_ringtones_df.merge( carveouts_df, how='left', on='release_id' ) merged_ringtones_df = merged_ringtones_df.merge( abacus_active_contracts_df, how='left', on=['track_id', 'track_isrc', 'upc', 'vendor_id'], suffixes=('', '_abacus'), ) merged_ringtones_df['mechanical_type'] = merged_ringtones_df['mechanical_type'].apply( parse_mechanical_type_to_list ) merged_ringtones_df['has_valid_oa_contract'] = ( (merged_ringtones_df['digital_track']) ) merged_ringtones_df['has_valid_abacus_contract'] = merged_ringtones_df['mechanical_type'].map( lambda m: 'digital' in m ) merged_ringtones_df = merged_ringtones_df[ (merged_ringtones_df['has_valid_oa_contract'] | merged_ringtones_df['has_valid_abacus_contract']) & ( ~merged_ringtones_df['has_us_carveout'] ) ] if merged_ringtones_df.empty: return pd.DataFrame(columns=['track_id', 'hfa_configuration_code']) merged_ringtones_df = merged_ringtones_df.reset_index().rename(columns={'index': '_original_order'}) filtered_ringtones_df = merged_ringtones_df.sort_values(by='_original_order').drop( columns=['_original_order'] ) filtered_ringtones_df = filtered_ringtones_df[filtered_ringtones_df['track_id'].notnull()] filtered_ringtones_df = filtered_ringtones_df[['track_id', 'hfa_configuration_code']] return filtered_ringtones_df def filter_abacus_contracts_by_priority( abacus_active_contract_df: pd.DataFrame, raw_digital_physical_df: pd.DataFrame, raw_ringtones_df: pd.DataFrame, ) -> pd.DataFrame: """Filter Abacus active contracts based on priority match. 1. ISRC match for term_type 'track' 2. UPC match for term_type 'product' 3. Vendor ID match for term_type 'label' Args: abacus_active_contract_df (pd.DataFrame): DataFrame of active Abacus contracts. raw_digital_physical_df (pd.DataFrame): Source digital/physical track data. raw_ringtones_df (pd.DataFrame): Source ringtone track data. Returns: pd.DataFrame: Filtered DataFrame with contract matches and metadata. """ merged_input_tracks_df = dataframes.merge_dataframes( [raw_digital_physical_df, raw_ringtones_df], columns=['track_id', 'track_isrc', 'upc', 'vendor_id'], drop_duplicates=True, ) matched_contracts = [] for _, track in merged_input_tracks_df.iterrows(): track_isrc = str(track['track_isrc']) track_id = str(track['track_id']) upc = str(track['upc']) vendor_id = track['vendor_id'] matched_contract = None matched_term_type = None # Priority 1: ISRC match where term_type is 'track' for _, contract in abacus_active_contract_df[ (abacus_active_contract_df['term_type'] == 'track') & (abacus_active_contract_df['account_id'] == vendor_id) ].iterrows(): if track_isrc in contract['attachments']: matched_contract = contract matched_term_type = 'track' break # Priority 2: UPC match where term_type is 'product' if matched_contract is None: for _, contract in abacus_active_contract_df[( (abacus_active_contract_df['term_type'] == 'product') & ( abacus_active_contract_df['account_id'] == vendor_id ) )].iterrows(): if upc in contract['attachments']: matched_contract = contract matched_term_type = 'product' break # Priority 3: Vendor ID match where term_type is 'label' if matched_contract is None: for _, contract in abacus_active_contract_df[( (abacus_active_contract_df['term_type'] == 'label') & ( abacus_active_contract_df['account_id'] == vendor_id ) )].iterrows(): if str(vendor_id) in contract['attachments']: matched_contract = contract matched_term_type = 'label' break if matched_contract is not None: matched_contracts.append( { 'track_id': int(track_id), 'track_isrc': track_isrc, 'upc': int(upc), 'vendor_id': vendor_id, 'matched_term_type': matched_term_type, 'mechanical_type': matched_contract['mechanical_type'], 'contract_id': matched_contract['contract_id'], } ) return pd.DataFrame( matched_contracts, columns=[ 'track_id', 'track_isrc', 'upc', 'vendor_id', 'matched_term_type', 'mechanical_type', 'contract_id', ], ) async def fetch_oa_active_contracts(vendor_ids: List[int]) -> pd.DataFrame: """Fetch OA active contracts for a list of vendor IDs asynchronously. Args: vendor_ids (List[int]): List of vendor IDs. Returns: Dict[str, Any]: Contracts data. """ try: contracts = await ows_contracts.get_oa_active_contracts(vendor_ids) contracts_df = dataframes.dict_to_dataframe( [ { 'vendor_id': int(vendor_id), 'physical_track': bool(contract.get('mechadmin_physical', False)), 'digital_track': bool(contract.get('mechadmin_digital', False)), } for vendor_id, contract in contracts.items() ], dtypes={ 'vendor_id': 'int64', 'physical_track': 'bool', 'digital_track': 'bool', }, ) return contracts_df except Exception as e: logger.exception(f'Unexpected Error while fetching OA Active Contracts: {e}') raise async def fetch_carveouts(release_ids: List[int]) -> pd.DataFrame: """Fetch carveouts data for the given list of release ids asynchronously. Args: release_ids (List[int]): List of release_ids to fetch carveouts Returns: Dict[str, Any]: Carveouts dataframe """ try: carveouts = await ows_carveouts_python.get_carveouts(release_ids) carveouts_df = dataframes.dict_to_dataframe( [ { 'release_id': int(release_id), 'has_us_carveout': any( 'US' in data.get(carveout_level, {}).get('country', []) for carveout_level in ('account', 'subaccount', 'product') ), } for release_id, data in carveouts.items() ], dtypes={ 'release_id': 'int64', 'has_us_carveout': 'bool', }, ) return carveouts_df except Exception as e: logger.exception(f'Unexpected error while fetching carveouts data: {e}') raise async def get_pending_hfa_request() -> List[Dict[str, str | int]]: """Fetch pending HFA requests using multiple microservices. Returns: List[Dict[str, str | int]]: A list of dictionaries containing trackId and hfaConfigurationCode. """ logger.info('Fetching pending HFA tracks') try: eligible_tracks = ows_product.get_hfa_eligible_tracks() if not eligible_tracks: return eligible_tracks raw_digital_physical_df = pd.DataFrame(eligible_tracks).query('hfa_configuration_code != "RT"') raw_ringtones_df = pd.DataFrame(eligible_tracks).query('hfa_configuration_code == "RT"') abacus_active_contract_df = fetch_and_load_abacus_active_contracts() filtered_abacus_active_contracts_df = filter_abacus_contracts_by_priority( abacus_active_contract_df, raw_digital_physical_df, raw_ringtones_df ) merged_tracks_df = dataframes.merge_dataframes( [raw_digital_physical_df, raw_ringtones_df], columns=['release_id', 'vendor_id'], drop_duplicates=True ) release_ids = merged_tracks_df['release_id'].dropna().unique().tolist() vendor_ids = merged_tracks_df['vendor_id'].dropna().unique().tolist() oa_active_contracts_df, carveouts_df = await asyncio.gather( fetch_oa_active_contracts(vendor_ids), fetch_carveouts(release_ids) ) eligible_digital_physical_tracks_df = filter_eligible_hfa_digital_physical_tracks( raw_digital_physical_df, oa_active_contracts_df, filtered_abacus_active_contracts_df, carveouts_df ) eligible_ringtones_df = filter_eligible_hfa_ringtones( raw_ringtones_df, oa_active_contracts_df, filtered_abacus_active_contracts_df, carveouts_df ) hfa_pending_tracks_df = dataframes.merge_dataframes( [eligible_digital_physical_tracks_df, eligible_ringtones_df], columns=['track_id', 'hfa_configuration_code'], drop_duplicates=False ) hfa_pending_tracks_df['track_id'] = hfa_pending_tracks_df['track_id'].astype(int) pending_tracks = [ { constants.FIELD_TRACK_ID: row['track_id'], constants.FIELD_HFA_CONF_CODE: row['hfa_configuration_code'], } for row in hfa_pending_tracks_df.to_dict(orient='records') ] if not pending_tracks: logger.warning('No pending HFA requests found!') return pending_tracks except Exception as ex: logger.exception(f'Unexpected error in get_pending_hfa_request: {ex}') raise