"""Used in UPC remapping.""" from ddex_ingester_common.constants.catalog_ingestion import \ DDEX_PROVIDER_TO_SOURCE_ID from ddex_ingester_common.constants.format_mapping import FORMAT_MAPPING from ddex_ingester_common.helpers.rds import run_rds_query from ddex_ingester_common.models.s3.body import Body as S3Body from ddex_ingester_common.models.state_machine.body import \ Body as StateMachineBody import config from constants import graphql_queries from constants.sql_queries import SELECT_REMAPPED_UPC, UPDATE_REMAPPED_UPC from constants.vendor import VENDOR_VALID_ARTIST_ID def get_project_id(logger, s3_context: S3Body) -> int: """Retrieve project id if the project exists, if not returns None.""" logger.info('Checking if project exists') # Project GraphQL queries need a null subaccount id to be 0 subaccount_id = s3_context.product.subaccount_id or 0 payload = { 'projectCode': s3_context.project.project_code, 'accountId': s3_context.product.vendor_id, 'subaccountId': subaccount_id } result = config.graphql_gateway.execute( graphql_queries.GET_PROJECT_BY_PROJECT_CODE, payload )['data']['projectByProjectCode'] logger.info(f'Ran GraphQL with payload {payload} and got result {result}') return result.get('projectId') if result else None def create_placeholder_project(logger, s3_context: S3Body, video=False) -> int: """Create a placeholder project and return the project id.""" logger.info('Creating placeholder project') artist_id = 1 vendor_id = s3_context.product.vendor_id if vendor_id in VENDOR_VALID_ARTIST_ID: artist_id = VENDOR_VALID_ARTIST_ID[vendor_id] elif video: raise Exception('Placeholder video project needs a valid artist id.') # Project GraphQL queries need a null subaccount id to be 0 subaccount_id = s3_context.product.subaccount_id or 0 payload = { 'data': { 'projectCode': s3_context.project.project_code, 'name': f'{s3_context.project.project_code} Remap Dummy', 'artistId': artist_id, 'accountId': vendor_id, 'subaccountId': subaccount_id } } result = config.graphql_gateway.execute( graphql_queries.CREATE_PROJECT, payload )['data']['createProject'] logger.info(f'Ran GraphQL with payload {payload} and got result {result}') return result['projectId'] def create_placeholder_audio_product( logger, context: StateMachineBody, project_id: int) -> str: """Create a placeholder product and return the new UPC for remapping.""" release_type = FORMAT_MAPPING.get(context.product.release_type) payload = { 'data': { 'projectId': project_id, 'productName': f'{context.product.upc} Remap Dummy', 'productHighlights': 'abc', 'format': release_type, 'accountId': context.product.vendor_id, 'subaccountId': context.product.subaccount_id, } } logger.info(f'Creating placeholder audio product with payload: {payload}') result = config.graphql_gateway.execute( graphql_queries.CREATE_AUDIO_PRODUCT, payload )['data']['createProduct'] return result['upc'] def create_placeholder_video_product( logger, s3_context: S3Body, project_id: int) -> str: """Create a placeholder product and return the new UPC for remapping.""" payload = { 'data': { 'create': { 'accountId': s3_context.product.vendor_id, 'subaccountId': s3_context.product.subaccount_id, 'projectId': str(project_id), 'typeOfVideo': s3_context.video.video_type, 'isrc': s3_context.video.isrc, } } } logger.info(f'Creating placeholder video product with payload: {payload}') result = config.graphql_gateway.execute( graphql_queries.CREATE_VIDEO_PRODUCT, payload )['data']['saveVideoSingleProduct'] return result['upc'] def get_remapped_upc_rows(logger, context: StateMachineBody) -> list: """Get the UPC remap rows.""" logger.info('Reading UPC remapping from RDS.') original_upc = context.product.upc source_id = DDEX_PROVIDER_TO_SOURCE_ID.get(context.ddex_provider) query_args = (original_upc, source_id) return run_rds_query( logger, config.RDS_HOST, config.RDS_DB_NAME, config.RDS_RW_USER, config.secrets_manager_client.get_cred('rds_read_write_password'), SELECT_REMAPPED_UPC, query_args, ) def update_remapped_upc(logger, context: StateMachineBody, new_upc: str): """Update a remapped UPC in RDS.""" original_upc = context.product.upc source_id = DDEX_PROVIDER_TO_SOURCE_ID.get(context.ddex_provider) query_args = (new_upc, original_upc, source_id) run_rds_query( logger, config.RDS_HOST, config.RDS_DB_NAME, config.RDS_RW_USER, config.secrets_manager_client.get_cred('rds_read_write_password'), UPDATE_REMAPPED_UPC, query_args, )