"""Interface for work with events.""" import datetime import json import logging from confluent_kafka import error as kafka_error, Producer from payee.config import Config from payee.connectors import sentry from payee.constants.constants import KAFKA_TOPICS from payee.utils.exception import KafkaSendEventException logger = logging.getLogger('kafka_producer_logger') class AccountPayeeProducer: """Account payee event producer.""" def __init__(self, topic): """Initialize AccountPayeeProducer.""" logger.info('Initializing AccountPayeeProducer') self._topic = topic producer_conf = { 'bootstrap.servers': Config.KAFKA_BOOTSTRAP_SERVERS, 'security.protocol': 'SSL', 'debug': 'all', } producer_logger = logging.getLogger('producer') producer_logger.setLevel(logging.DEBUG) self._producer = Producer(producer_conf, logger=producer_logger) logger.info('AccountPayeeProducer is initialized.') def __enter__(self): """Enter producer context.""" logger.info('Entering producer context.') return self def __exit__(self, exc_type, exc_val, exc_tb): """Exit producer context.""" logger.info('Exiting producer context.') if exc_type: logger.error(f'Context manager failure: {exc_type}, {exc_val}, {exc_tb}') self._producer.flush() logger.info('Messages successfully flushed.') def produce_event(self, event_key, event_payload): """Produce new account payee event. Args: event_key (str): event key (account_payee_id as a part). event_payload (dict): event payload. """ try: self._producer.produce( topic=self._topic, key=event_key, value=json.dumps(event_payload).encode('utf-8'), on_delivery=_default_on_delivery, ) except kafka_error.KeySerializationError as e: sentry.send_to_sentry( f'Could not serialize key. Error: {str(e)}', str(e), 'error', f'Error! Code: {type(e).__name__}, Message, {str(e)}', ) raise KafkaSendEventException(f'Could not serialize key. Error: {str(e)}') except kafka_error.ValueSerializationError as e: sentry.send_to_sentry( f'Could not serialize value. Error: {str(e)}', str(e), 'error', f'Error! Code: {type(e).__name__}, Message, {str(e)}', ) raise KafkaSendEventException(f'Could not serialize value. Error: {str(e)}') except kafka_error.KafkaException as e: sentry.send_to_sentry( f'Kafka message send failed. Error: {str(e)}', str(e), 'error', f'Error! Code: {type(e).__name__}, Message, {str(e)}', ) raise KafkaSendEventException(f'Kafka message send failed. Error: {str(e)}') def emit_payee_event(account_payee_id: int, action_name: str): """Emit payee event from passed data. Args: account_payee_id (int): Payee id. action_name (str): payment/tax eligibility. """ event_key, event_payload = _create_payee_event_key_and_payload(account_payee_id) if action_name not in Config.KAFKA_PRODUCERS_BY_ACTION_NAME: Config.KAFKA_PRODUCERS_BY_ACTION_NAME[action_name] = AccountPayeeProducer( KAFKA_TOPICS[action_name], ) with Config.KAFKA_PRODUCERS_BY_ACTION_NAME[action_name] as producer: producer.produce_event(event_key, event_payload) def _default_on_delivery(err, msg): if err is not None: sentry.send_to_sentry( f'Delivery failed for payee `{msg.key()}`: {str(err)}', err, 'error', f'Error! Code: {type(err).__name__}, Message, {str(err)}', ) raise KafkaSendEventException( f'Delivery failed for payee `{msg.key()}`: {str(err)}' ) elif msg: sentry.send_to_sentry( 'Payee record {} successfully produced to {} [{}] at offset {}'.format( msg.key(), msg.topic(), msg.partition(), msg.offset() ), {}, 'info', msg, ) def _create_payee_event_key_and_payload(account_payee_id): """Compose event key and payload.""" key = f'account_payee_id_{account_payee_id}' payload = { 'timestamp': datetime.datetime.now().isoformat(), 'accountPayee': {'accountPayeeId': f'{account_payee_id}'}, } return key, payload