"""Lambda function module for create_product.""" from typing import Dict from common.connectors.mysql_connector import execute_mysql_query from common.constants.product_nfd import SME_ANALYTICS_DUMMY from common.constants.product_status import IN_CONTENT from common.constants.role_mappings import PRODUCT_ROLES from common.models.state_machine.grps_ingestion_context import \ GrpsIngestionContext from common.schemas.state_machine_schema import ProductType, StateMachineSchema from lambdacommon.graphql import graphql import config from config import graphql_gateway from src.constants import sql_queries from src.constants.queries import CREATE_AUDIO_PRODUCT, \ CREATE_OR_UPDATE_VIDEO_PRODUCT, GET_PRODUCT_BY_UPC, UPDATE_PRODUCT from src.exceptions import CreateProductException, \ InContentProductIsNotSmeAnalyticsDummyException logger = config.app_logger def handler(event, context): """Create product handler. Checks if a product exists by UPC (UPC is the product identifier). If it does, returns existing product data. If not, creates a new product using context data. """ logger.info(f'Triggered create_product: {event}') sm_context = StateMachineSchema().load(event) correlation_id = sm_context.correlation_id graphql_gateway.set_headers( { 'Orchard-User-Id': config.OA_USER, 'Correlation-Id': correlation_id, } ) try: graphql_result = check_for_product(sm_context) if graphql_result: if (graphql_result['status'] == IN_CONTENT and # noqa graphql_result['notForDistribution'] != SME_ANALYTICS_DUMMY and graphql_result['deletions'] == 'N'): # noqa raise InContentProductIsNotSmeAnalyticsDummyException() sm_context.product.product_id = graphql_result['productId'] logger.info(f'Product already exists with upc: ' f'{sm_context.product.upc} and product_id: ' f'{graphql_result["productId"]}.Do update.') sm_context.product.status = graphql_result['status'] update_product(sm_context) else: create_product(sm_context) except graphql.GraphQLError as err: message = str(err) logger.error(f'GraphQL error in create_product: {message}') raise CreateProductException(message) from err return StateMachineSchema().dump(sm_context) def create_product(sm_context: StateMachineSchema): """Create video/audio products.""" if sm_context.product_type == ProductType.AUDIO: result = create_audio_product(sm_context) sm_context.product.product_id = result['productId'] # releaseDate/selestartDate params are not supported in creation update_product(sm_context) elif sm_context.product_type == ProductType.VIDEO: result = create_video_product(sm_context) sm_context.product.product_id = result['productId'] # ows-video rejects release dates in past update_product(sm_context) def check_for_product(context: GrpsIngestionContext) -> Dict: """Check for existence of a product by UPC. UPC serves as the product identifier. Args: context: State machine context containing product info. Returns: dict: Product data if found, empty dict otherwise. """ upc = context.product.upc if not upc: logger.info('No UPC found, skipped check_for_product') return {} logger.info(f'Checking for existing product with upc: {upc}') payload = { 'upc': upc } result = graphql_gateway.execute( GET_PRODUCT_BY_UPC, payload )['data']['productByUpc'] if result: logger.info( f'Found existing product for upc={upc}') return result def create_audio_product(context: GrpsIngestionContext) -> Dict: """Create a product from context data. Uses product, project, and label_participants information from the state machine context (sourced from GRPS). UPC serves as the product identifier. Args: context: State machine context containing all product info. Returns: dict: Created product data from GraphQL response. """ participations = build_participations(context) payload = { 'data': { 'productName': context.product.product_name, 'productHighlights': 'abc', 'projectId': context.project.project_id, 'accountId': context.product.vendor_id, 'subaccountId': context.product.subaccount_id, 'upc': context.product.upc, 'notForDistribution': context.product.not_for_distribution, 'participations': participations, 'manufacturerUpc': context.product.upc, 'format': context.product.release_type, 'imprint': context.product.imprint, 'productCode': context.product.catalog_number, } } logger.info(f'Running create product with payload:\n{payload}') result = graphql_gateway.execute( CREATE_AUDIO_PRODUCT, payload )['data']['createProduct'] if result: logger.info( f'Product created for upc: {context.product.upc}') logger.info(result) if context.placeholder_upc_ingestion: product_id = result['productId'] context.product.upc = result['upc'] update_display_upc(product_id, context.product.display_upc) store_placeholder_upc(context, context.product.display_upc, context.product.upc) return result def build_participations(context: GrpsIngestionContext) -> list: """Build participations payload from context data. Matches product display_artists to label_participants by name to retrieve label_participant_uuid for the GraphQL mutation. Args: context: State machine context containing artist and label_participant info. Returns: list: List of participation dicts for GraphQL payload. """ if not context.product.display_artists: return [] label_participants = context.label_participants or [] participations = [] for artist in context.product.display_artists: label_participant = _find_label_participant( label_participants, artist.name ) for role in (artist.roles or []): if role not in PRODUCT_ROLES: continue participation = { 'labelParticipantUuid': ( label_participant.label_participant_uuid if label_participant else None ), 'role': role, } participations.append(participation) return participations def _find_label_participant(label_participants, artist_name): """Find a label participant by artist name. Args: label_participants: List of LabelParticipant objects. artist_name: Name to match. Returns: LabelParticipant or None. """ for lp in label_participants: if lp.name == artist_name: return lp return None def update_product(context: GrpsIngestionContext) -> Dict: """Update product.""" participations = build_participations(context) payload = { 'data': { 'productId': context.product.product_id, 'releaseDate': context.product.release_date, 'saleStartDate': context.product.sale_start_date, 'productName': context.product.product_name, 'participations': participations, 'manufacturerUpc': context.product.upc, 'format': context.product.release_type, 'imprint': context.product.imprint, 'productCode': context.product.catalog_number } } logger.info(f'Updating product release date with payload: {payload}') result = graphql_gateway.execute( UPDATE_PRODUCT, payload ) if result: logger.info( f'Product release data updated: {context.product.product_id}') logger.info(result) return result def store_placeholder_upc(context, upc, placeholder_upc): """Store upc to placeholder upc mapping.""" query_args = ( upc, placeholder_upc, context.product.grid, f'grps_ingestor run {context.grps_ingestion_id}' ) execute_mysql_query( logger, config.RDS_HOST, config.RDS_DB_NAME, config.RDS_RW_USER, config.RDS_PASSWORD, sql_queries.INSERT_PLACEHOLDER_UPC, query_args, ) def update_display_upc(product_id, display_upc): """Update products display_upc. Unsafe! Use only for ddex product creations with placeholder_upc approach """ logger.info( f'Updating display upc to {display_upc} for product {product_id}') update_upc_response = config. \ ows_client.put('ows-product-digital', path=f'/product/{product_id}/display_upc', json={'display_upc': display_upc}) if update_upc_response.status_code != 200: raise Exception('Unable to update display_upc for the product' f'error message : {update_upc_response.text}') def create_video_product(context: GrpsIngestionContext): """Create video product.""" payload = { 'create': { 'upc': context.product.upc, 'accountId': context.product.vendor_id, 'subaccountId': context.product.subaccount_id, 'projectId': str(context.project.project_id), 'typeOfVideo': 'Other', 'imprint': context.product.imprint, 'productCode': context.product.catalog_number, 'notForDistribution': context.product.not_for_distribution, 'videoTitle': context.product.product_name, } } logger.info(f'Creating video product with payload: {payload}') result = graphql_gateway.execute( CREATE_OR_UPDATE_VIDEO_PRODUCT, {'data': payload} )['data']['saveVideoSingleProduct'] return result