"""replace_with_placeholder_upc.""" import json import uuid import boto3 import config from constants import graphql_queries, sql_queries from ddex_ingester_common.constants.status import IN_CONTENT from ddex_ingester_common.constants.swb_deal_types import \ SME_ANALYTICS_NOT_FOR_DISTRIBUTION from ddex_ingester_common.helpers.rds import run_rds_query from ddex_ingester_common.helpers.s3_ddex import load_ddex_json from ddex_ingester_common.logging import utils as logging_utils from ddex_ingester_common.models.s3.body import Body as S3Body from ddex_ingester_common.models.state_machine.body import ( Body as StateMachineBody) from ddex_ingester_common.schemas.s3_schema import S3Schema from ddex_ingester_common.schemas.state_machine_schema import \ StateMachineSchema logger = logging_utils.get_logger(config.app_logger) def handler(event, context): """replace_with_placeholder_upc handler.""" logger.info(f'Triggered replace_with_placeholder_upc: {event}') context = StateMachineSchema().load(event) s3_context = S3Schema().load(load_ddex_json(event)) correlation_id = context.correlation_id or str(uuid.uuid4()) context.correlation_id = correlation_id s3_context.correlation_id = correlation_id logging_utils.update_logger_correlation_id(logger, correlation_id) logging_utils.update_logger_with_message_ids( logger, context.message_id, context.message_thread_id, context.execution_name ) config.graphql_gateway.set_headers( { 'Orchard-User-Id': config.OA_USER, 'Correlation-Id': correlation_id, } ) upc = s3_context.product.upc placeholder_upc = get_placeholder_upc(upc) if placeholder_upc: logger.info(f'Found placeholder_upc {placeholder_upc} for upc: {upc}') display_upc = upc upc = placeholder_upc else: product = get_product_by_upc(upc) logger.info(f'Get product by upc :{product} , upc : {upc}') if product and product[ 'notForDistribution'] != SME_ANALYTICS_NOT_FOR_DISTRIBUTION and \ product['status'] != IN_CONTENT: logger.info( f'Product with UPC {upc} already exists, ' f'With notForDistribution : {product["notForDistribution"]} ' f'Release status is {product["status"]} ' 'Will insert product with new upc as placeholder') display_upc = upc upc = None else: return StateMachineSchema().dump(context) context.product.upc = upc context.product.display_upc = display_upc s3_context.product.upc = upc s3_context.product.display_upc = display_upc save_s3_context(context, s3_context) return StateMachineSchema().dump(context) def save_s3_context(context: StateMachineBody, s3_context: S3Body): """Save s3 context.""" s3_client = boto3.client('s3') json_file_path = f'{context.key}parsed_ddex.json' s3_client.put_object( Bucket=context.bucket, Key=json_file_path, Body=json.dumps(S3Schema().dump(s3_context)).encode(encoding='UTF-8') ) def get_placeholder_upc(upc): """Get placeholder upc if exists .""" query_args = ( upc ) result = run_rds_query( logger, config.RDS_HOST, config.RDS_DB_NAME, config.RDS_RW_USER, config.RDS_PASSWORD, sql_queries.GET_PLACEHOLDER_UPC, query_args, ) if result: return result[0].get('placeholder_upc') else: return None def get_product_by_upc(upc) -> dict: """Check if product already exists .""" graphql_result = config.graphql_gateway.execute( graphql_queries.GET_PRODUCT_BY_UPC, {'upc': upc}) return graphql_result['data']['productByUpc']