"""test_rating_computation test module.""" import base64 import hashlib import json import operator import os from datetime import datetime from datetime import timedelta from kafka import KafkaConsumer from kafka import KafkaProducer import pytest import requests CLIENT_ID = os.environ.get('CLIENT_ID') BOOTSTRAP_SERVERS = os.environ.get('BOOTSTRAP_SERVERS') MONITOR_IDS = [ 42903524, # qa-kafka-connect-sforce-sink-dlq-monitor 42903525, # qa-kafka-connect-sforce-sink-error-monitor ] DD_API_URL = 'https://api.datadoghq.com/api/v1/monitor' DD_HEADERS = { 'DD-API-KEY': os.environ.get('DD_API_KEY'), 'DD-APPLICATION-KEY': os.environ.get('DD_APPLICATION_KEY') } MONITOR_MUTE_DURATION = 15 # mute monitor for 15 minutes @pytest.fixture(scope='session', autouse=True) def mute_monitors(): """Mute DataDog monitors.""" mute_duration = datetime.now() + timedelta(minutes=MONITOR_MUTE_DURATION) for monitor_id in MONITOR_IDS: requests.post( f'{DD_API_URL}/{monitor_id}/mute', params={'end': mute_duration.timestamp()}, headers=DD_HEADERS) @pytest.fixture def email(request): """Generate random email.""" return 'test_' + base64.b64encode( os.urandom(6)).decode('ascii') + '@test.com' @pytest.fixture() def producer(request): """Create producer object.""" return KafkaProducer( bootstrap_servers=BOOTSTRAP_SERVERS, security_protocol='SSL', client_id=CLIENT_ID, linger_ms=0, api_version='2.6.0') def generate_salesforce_id(email: str, artist_name: str) -> bytes: """Generate salesforce id hash.""" data = ''.join(s.lower() for s in [email, artist_name]) data += 'mdtSXLGxWbQWTmb5lFQZNvS4iAKFjbJT8HJOooA' email_hash = hashlib.sha256(data.encode()) return email_hash.hexdigest().encode() def create_consumer(event, group_id): """Create consumer object.""" return KafkaConsumer(event, client_id=CLIENT_ID, bootstrap_servers=BOOTSTRAP_SERVERS, security_protocol='SSL', auto_offset_reset='latest', enable_auto_commit=False, group_id=group_id, consumer_timeout_ms=3000 ) def get_last_partition_position(consumer): """Return partition and latest partition position.""" consumer.poll(500) consumer.seek_to_end() topic_partition = consumer.assignment() partition = list(topic_partition)[0] return partition, consumer.position(partition) def configure_data(email, file): """Configure variable data from test data.""" data = json.load(file) data['Email'] = email object_id = generate_salesforce_id(data['Email'], data['Company']) data['KafkaMessageHeaders__c']['CamelHeader.sObjectIdValue'] = \ object_id.decode() return data, object_id @pytest.mark.parametrize('test_data,score,sign', [('approve_data.json', 'Review', '>='), ('review_data.json', 'Review', '<'), ('blacklist_review_data.json', 'Review', None), ]) def test_rating_computation(email, producer, test_data, score, sign): """Test correct score returned for approve and review scenarios.""" consumer = create_consumer('event.gdaSignup', 'integration-test') partition, start_position = get_last_partition_position(consumer) f = open(f'tests/test_data/{test_data}') data, object_id = configure_data(email, f) producer.send( 'event.gdaSubmitApplication', json.dumps(data).encode(), headers=[ ('OrchardHeader.CorrelationId', b'uuid-uuid'), ('CamelHeader.sObjectIdValue', object_id) ] ) producer.close() message_found = False consumer.poll(30000) consumer.seek(partition, start_position) for index, message in enumerate(consumer): message_found = True assert index < 1 message_value = json.loads(message.value.decode('utf-8')) assert message_value['Email'] == data['Email'] assert message_value['LeadScore__c'] == score if score in ['Approve', 'Review'] and sign is not None: operators = {'>=': operator.ge, '<': operator.lt} op_func = operators[sign] assert op_func(message_value['SpotifyFollowers__c'], 1000) assert op_func(message_value['SpotifyMonthlyListeners__c'], 3000) consumer.close() assert message_found def test_sales_force_api_failure(email, producer): """Test sales force api failure.""" consumer = create_consumer('dlq.gdaSignup', 'integration-test') partition, start_position = get_last_partition_position(consumer) f = open('tests/test_data/approve_data.json') data, object_id = configure_data(email, f) data['RecordTypeId'] = 'junk' producer.send( 'event.gdaSubmitApplication', json.dumps(data).encode(), headers=[ ('OrchardHeader.CorrelationId', b'uuid-uuid'), ('CamelHeader.sObjectIdValue', object_id) ] ) producer.close() message_found = False consumer.poll(30000) consumer.seek(partition, start_position) for index, message in enumerate(consumer): message_found = True assert index < 1 message_value = json.loads(message.value.decode('utf-8')) assert message_value['Email'] == data['Email'] consumer.close() assert message_found