"""Lambda gda_calculate_score function module.""" from kafka import KafkaProducer from lambdacommon.common_config import logger from lambdacommon.common_config import SENTRY_DSN from marshmallow.exceptions import ValidationError import sentry_sdk from sentry_sdk.integrations.aws_lambda import AwsLambdaIntegration from sentry_sdk import capture_exception import config from . import constants from .connectors import ows_blacklist_manager from .connectors import voucherify_manager from .schema import CalculateScoreSchema sentry_sdk.init( dsn=SENTRY_DSN, environment=config.ENVIRONMENT, integrations=[AwsLambdaIntegration(timeout_warning=True)] ) def _send_message(event: dict, salesforce_id: str, correlation_id: str): """Send message to an MSK topic.""" producer = KafkaProducer(**{ 'bootstrap_servers': config.BOOTSTRAP_SERVERS, 'security_protocol': 'SSL', 'client_id': config.CLIENT_ID, 'linger_ms': config.LINGER_MS, 'api_version': config.KAFKA_API_VERSION }) validation_errors = CalculateScoreSchema().validate(data=event) if validation_errors: logger.error( f'Schema validation in gda_calculate_score lambda has failed' f'with the following errors: {str(validation_errors)}') raise ValidationError(str(validation_errors)) producer.send( config.TOPIC_NAME, CalculateScoreSchema().dumps(event).encode(), headers=[ (constants.ID_HEADER, salesforce_id.encode()), (constants.CORRELATION_ID_HEADER, correlation_id.encode()), ]) producer.flush() def _calculate_score(**kwargs: dict) -> dict: lead_score = {'score': None, 'message': None} spotify_followers = kwargs.get('spotify_followers') or 0 spotify_monthly_listeners = kwargs.get('spotify_monthly_listeners') or 0 artist_or_band = kwargs.get('artist_or_band') or '' first_name = kwargs.get('first_name') or '' last_name = kwargs.get('last_name') or '' existing_label_participants = kwargs.get('label_participants') data = ' '.join([artist_or_band, first_name, last_name]) country = kwargs.get('country') or '' comment = kwargs.get('comment') or '' # Checking Voucherify voucher and validate if exists voucher_code = kwargs.get('voucher_code') orchard_lead_id = kwargs.get('orchard_lead_id') valid_voucher = voucherify_manager.validate_voucher(voucher_code, orchard_lead_id) if valid_voucher and valid_voucher.get('valid'): lead_score['valid_voucher'] = valid_voucher lead_score['score'] = constants.APPROVE_STATUS_VALUE lead_score['message'] = constants.COMMENTS_VOUCHER_APPROVED_ARTIST response = ows_blacklist_manager.validate_blacklist(data) if country == 'RUS' or country == 'TUR': lead_score['score'] = constants.REJECT_STATUS_VALUE lead_score['message'] = 'Access has been rejected.' return lead_score if response.status == 400: logger.error( 'Blacklist validation failed. Data: [%s], Errors: [%s]', data, response.errors) lead_score['score'] = constants.REVIEW_STATUS_VALUE lead_score['message'] = constants.COMMENTS_BLOCKED_ARTIST elif response.status == 500: logger.error(response.message) lead_score['score'] = constants.REVIEW_STATUS_VALUE elif existing_label_participants and not valid_voucher.get('valid'): lead_score['score'] = constants.REVIEW_STATUS_VALUE lead_score['message'] = constants.COMMENTS_EXISTING_LABEL elif comment != '': lead_score['score'] = constants.REVIEW_STATUS_VALUE lead_score['message'] = constants.COMMENT_EXISTED if lead_score.get('score'): return lead_score min_followers = constants.SPOTIFY_FOLLOWERS_APPROVE_THRESHOLD min_listeners = constants.SPOTIFY_MONTHLY_LISTENERS_APPROVE_THRESHOLD upper_bound_followers = constants.SPOTIFY_FOLLOWERS_UPPER_THRESHOLD upper_bound_listeners = constants.SPOTIFY_MONTHLY_LISTENERS_UPPER_THRESHOLD if (min_listeners <= spotify_monthly_listeners <= upper_bound_listeners) and ( min_followers <= spotify_followers <= upper_bound_followers) and ( country in constants.COUNTRIES_WITHOUT_REVIEW) and comment == '': lead_score['score'] = constants.REVIEW_STATUS_VALUE lead_score['message'] = constants.COMMENTS_NON_REFERRAL else: lead_score['score'] = constants.REVIEW_STATUS_VALUE if (spotify_monthly_listeners > upper_bound_listeners) or ( spotify_followers > upper_bound_followers): lead_score['message'] = constants.COMMENTS_EXCEEDED_UPPER_BOUND return lead_score def _merge_dicts_from_parallel_branches(event: list) -> dict: """Merge dicts from list of dicts preferring values over nulls.""" merged_event = {} for d in event: for key, value in d.items(): if not merged_event.get(key): merged_event[key] = value else: merged_event[key] = merged_event.get(key) or value return merged_event def handler(event: list, context: dict) -> dict: """Lambda entry point.""" try: event = _merge_dicts_from_parallel_branches(event) key_to_field_name_mapping = { 'label_participants': 'ExistingLabelParticipants__c', 'deezer_url': 'DeezerUrl__c', 'deezer_followers': 'DeezerFans__c', 'facebook_url': 'FacebookUrl__c', 'facebook_followers': 'FacebookFollowers__c', 'instagram_url': 'InstagramUrl__c', 'instagram_followers': 'InstagramFollowers__c', 'soundcloud_url': 'SoundcloudUrl__c', 'soundcloud_followers': 'SoundcloudFollowers__c', 'spotify_url': 'SpotifyUrl__c', 'spotify_followers': 'SpotifyFollowers__c', 'spotify_monthly_listeners': 'SpotifyMonthlyListeners__c', 'tiktok_url': 'TiktokUrl__c', 'tiktok_followers': 'TiktokFollowers__c', 'twitter_url': 'TwitterUrl__c', 'twitter_followers': 'TwitterFollowers__c', 'youtube_url': 'YoutubeUrl__c', 'youtube_followers': 'YoutubeSubscribers__c' } values_from_gda_knowledge_graph_result = { key_to_field_name_mapping[key]: event.get( 'lambda_gda_knowledge_graph_result', {}).get(key) for key in key_to_field_name_mapping.keys() } salesforce_id = event.get('KafkaMessageHeaders__c').get(constants.ID_HEADER) correlation_id = event.get('KafkaMessageHeaders__c').get(constants.CORRELATION_ID_HEADER) logger.info( f'Sending message with Salesforce id {salesforce_id} ' f' and Correlation id {correlation_id} to MSK' f' from the gda_score_calculation lambda...') for field_name, value in values_from_gda_knowledge_graph_result.items(): # we always prefer user input to values we get from the knowledge graph event.update({field_name: event.get(field_name) or value}) # except of the results of the scrapers, we use them if spotify_monthly_listeners := event.get( 'lambda_gda_spotify_scraper_result', {}).get( 'spotify_monthly_listeners'): event.update({'SpotifyMonthlyListeners__c': spotify_monthly_listeners}) if youtube_subscribers := event.get( 'lambda_gda_youtube_scraper_result', {}).get( 'youtube_subscribers'): event.update({'YoutubeSubscribers__c': youtube_subscribers}) if tiktok_followers := event.get( 'lambda_gda_tiktok_scraper_result', {}).get( 'tiktok_followers'): event.update({'TiktokFollowers__c': tiktok_followers}) participants = event.get('ExistingLabelParticipants__c') if participants: event.update({'ExistingLabelParticipants__c': participants.split(',')[0]}) if first_name := event.get('FirstName'): event.update({'FirstName': first_name.replace('\\', '')}) if last_name := event.get('LastName'): event.update({'LastName': last_name.replace('\\', '')}) if company := event.get('Company'): event.update({'Company': company.replace('\\', '')}) if country := event.get('Country'): event.update({'Country': country.replace('\\', '')}) if comment := event.get('Comments__c'): event.update({'Comments__c': comment.replace('\\', '')}) lead_score = _calculate_score( spotify_monthly_listeners=event.get('SpotifyMonthlyListeners__c'), spotify_followers=event.get('SpotifyFollowers__c'), artist_or_band=event.get('Company'), first_name=event.get('FirstName'), last_name=event.get('LastName'), voucher_code=event.get('VoucherCode__c'), orchard_lead_id=salesforce_id, label_participants=participants, country=event.get('Country'), comment=event.get('Comments__c') ) if lead_score: # Checking for Voucher Code Owner if 'valid_voucher' in lead_score and \ 'voucher_code_owner' in lead_score['valid_voucher']: voucher_code_owner = lead_score['valid_voucher']['voucher_code_owner'] event.update({'VoucherCodeOwner__c': voucher_code_owner}) if 'consultant' in lead_score['valid_voucher']: event.update({'Consultant__c': lead_score['valid_voucher']['consultant']}) # Checking for message if 'message' in lead_score: event.update({'Comments__c': lead_score['message']}) event.update({'LeadScore__c': lead_score['score']}) # delete redundant (for Salesforce connector) fields event.pop('KafkaMessageHeaders__c', None) event.pop('lambda_gda_knowledge_graph_result', None) event.pop('lambda_gda_spotify_scraper_result', None) event.pop('lambda_gda_tiktok_scraper_result', None) event.pop('lambda_gda_youtube_scraper_result', None) event.pop('Error', None) event.pop('Cause', None) event.pop('_datadog', None) # we're sending an original event, so we send AmazonUrl__c, OtherUrl__c, # etc if provided in that original event _send_message(event, salesforce_id, correlation_id) lead_score_result = event.get('LeadScore__c') logger.info( f'Message with Salesforce id {salesforce_id}, LeadScore__c {lead_score_result} ' f' and Correlation id {correlation_id} sent to MSK' f' from the gda_score_calculation lambda') return event except Exception as e: capture_exception(e) logger.exception(str(e)) raise e