"""Lambda function module.""" import uuid import config from config import graphql_gateway from constants import queries from ddex_ingester_common.constants.ddex_providers import ( SME, SME_ANALYTICS_PROVIDER) from ddex_ingester_common.helpers.s3_ddex import load_ddex_json from ddex_ingester_common.lambda_exceptions import ( ArtistIdNotFoundException, DuplicateProjectCodeException, ProjectCodeMismatchException, SetProjectException ) 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 ) from lambdacommon.graphql import graphql from marshmallow.utils import get_value logger = logging_utils.get_logger(config.app_logger) def handler(event, context): """Set project handler.""" logger.info(f'Triggered set_project: {event}') context = StateMachineSchema().load(event) s3_data = S3Schema().load(load_ddex_json(event)) correlation_id = context.correlation_id or str(uuid.uuid4()) 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 ) graphql_gateway.set_headers( { 'Orchard-User-Id': config.OA_USER, 'Correlation-Id': correlation_id, } ) project_code_mismatch = check_project_code_mismatch( s3_data, context, context.product.upc ) try: graphql_result = check_for_project(context, s3_data) if not graphql_result: graphql_result = create_project(context, s3_data) # for SME (Switchboard and Som Livre) and SMEAnalyticsDummy # don't update project if code mismatch elif not project_code_mismatch: context.project_id = graphql_result.get('projectId') graphql_result = update_project(context, s3_data) except graphql.GraphQLError as err: if context.ddex_provider == SME: message = err.get_response_body().get('message') raise SetProjectException(message) from err raise return StateMachineSchema().dump(context) def check_project_code_mismatch( s3_ddex: S3Body, context: StateMachineBody, upc: str) -> bool or None: """Check if the UPC has a different project associated to it.""" if not upc: logger.info('No UPC found, skipped check_project_code_mismatch') return ddex_project_code = \ s3_ddex.project.project_code if s3_ddex.project else None graphql_result = graphql_gateway.execute( queries.GET_PRODUCT_BY_UPC, {'upc': upc} )['data']['productByUpc'] logger.info(f'Received the project for upc {upc} of {graphql_result}') context_vendor = context.product.vendor_id graphql_vendor = get_value( graphql_result, 'vendorId' ) or None graphql_project_code = get_value( graphql_result, 'project.projectCode' ) # Remove 'SONY:id:' prefix if present if ddex_project_code: ddex_project_code = ddex_project_code.replace('SONY:id:', '') if graphql_project_code: graphql_project_code = graphql_project_code.replace('SONY:id:', '') if graphql_project_code and graphql_project_code != ddex_project_code: message = (f'GraphQL project code "{graphql_project_code}" does ' f'not match DDEX project code "{ddex_project_code}" ' f'DDEX Vendor ID: {context_vendor} GraphQL ' f'Vendor ID: {graphql_vendor}') logger.error(message) if context.ddex_provider in (SME, SME_ANALYTICS_PROVIDER): # for SME (Switchboard and Som Livre) and SMEAnalyticsDummy # don't blow up if code mismatch return True raise ProjectCodeMismatchException(message) if graphql_vendor and int(graphql_vendor) != int(context_vendor): message = (f'UPC {upc} has Vendor ID: {graphql_vendor} which does not ' f'match the DDEX Vendor ID: {context_vendor}') logger.error(message) raise ProjectCodeMismatchException(message) def create_project(context, s3_data): """Create a project from given params.""" # Project GraphQL queries need a null subaccount id to be 0 subaccount_id = context.product.subaccount_id or 0 vendor_id = context.product.vendor_id """ This check is in place temporarily while extended metadata is not mandatory. When we are starting to see this included with all DDEX product then we will no longer need to fake the data. """ project_name = s3_data.project.name if s3_data.project else\ str(uuid.uuid4())[:10] project_code = s3_data.project.project_code if s3_data.project else\ str(uuid.uuid4())[:10] description = s3_data.project.description if s3_data.project else\ str(uuid.uuid4())[:18] logger.info( f'Running create project with vendor_id: {vendor_id}' f' and subaccount id: {subaccount_id}' f' and project code: {project_code}') payload = { 'data': { 'projectCode': project_code, 'name': project_name, 'artistId': retrieve_artist_id(s3_data), 'description': description, 'accountId': vendor_id, 'subaccountId': subaccount_id } } try: result = graphql_gateway.execute( queries.CREATE_PROJECT, payload )['data']['createProject'] except graphql.GraphQLError as err: # We need to check for this error because it means the project was # created by another execution after we checked if it existed. # This means if we run the lambda again, we will see the project exists # and the execution can continue without issue. # We raise this exception instead of the GraphQLError as Terraform # catches this particular exception and runs the lambda again. if f"Project code '{project_code}' already exists" in str(err): raise DuplicateProjectCodeException(err) raise if result: context.project_id = result.get('projectId') return result def check_for_project(context, s3_data): """Check for existence of a project.""" # Project GraphQL queries need a null subaccount id to be 0 subaccount_id = context.product.subaccount_id or 0 vendor_id = context.product.vendor_id project_code = s3_data.project.project_code if s3_data.project else\ str(uuid.uuid4())[:18] logger.info( f'Running get project with vendor_id: {vendor_id}' f' and subaccount id: {subaccount_id}' f' and project code: {project_code}') payload = { 'projectCode': project_code, 'accountId': vendor_id, 'subaccountId': subaccount_id } result = graphql_gateway.execute( queries.GET_PROJECT_BY_PROJECT_CODE, payload )['data']['projectByProjectCode'] logger.info(f'Get project response: {result}') return result def update_project(context, s3_data): """Update existing project.""" # Project GraphQL queries need a null subaccount id to be 0 subaccount_id = context.product.subaccount_id or 0 vendor_id = context.product.vendor_id project_id = context.project_id """ This check is in place temporarily while extended metadata is not mandatory. When we are starting to see this included with all DDEX product then we will no longer need to fake the data. """ project_name = s3_data.project.name if s3_data.project else\ str(uuid.uuid4())[:10] description = s3_data.project.description if s3_data.project else\ str(uuid.uuid4())[:10] project_code = s3_data.project.project_code if s3_data.project else\ None logger.info( f'Running update project with vendor_id: {vendor_id}' f' and subaccount id: {subaccount_id}' f' and project code: {project_code}') payload = { 'data': { 'name': project_name, 'artistId': retrieve_artist_id(s3_data), 'projectId': project_id } } if description: payload['data']['description'] = description logger.info(f'Payload : {payload}') result = graphql_gateway.execute( queries.UPDATE_PROJECT, payload )['data']['updateProject'] return result def retrieve_artist_id(s3_data): """Retrieve project artist id from artist list in S3 data.""" artist_id = None if s3_data.label_participants and s3_data.project.artist: artist = None for label_participant in s3_data.label_participants: if label_participant.name == s3_data.project.artist.name: artist = label_participant if artist: artist_id = int(artist.artist_id) if not artist_id: raise ArtistIdNotFoundException( f'Artist id not found for name: {s3_data.project.artist.name}' ) return artist_id