"""Interface for work with events.""" import datetime import json import logging from confluent_kafka import error as kafka_error from confluent_kafka import Producer from abacus_account.config import Config from abacus_account.connectors import sentry from abacus_account.constants.constants import ACCOUNT_KAFKA_EVENT_NAMES, KAFKA_TOPICS from abacus_account.utils.exception import KafkaSendEventException logger = logging.getLogger('kafka_producer_logger') class AccountEventProducer: """Account event producer.""" def __init__(self, topic): """Initialize AccountInfoProducer.""" logger.info('Initializing AccountEventProducer') self._topic = topic producer_conf = { 'bootstrap.servers': Config.KAFKA_BOOTSTRAP_SERVERS, 'security.protocol': 'SSL', 'debug': 'all' } self._producer = Producer(producer_conf) logger.info('AccountEventProducer 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: print('Trying to produce event...') self._producer.produce( topic=self._topic, key=event_key, value=json.dumps(event_payload).encode('utf-8'), on_delivery=_default_on_delivery) print('Event was successfully produced.') 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_account_tax_id_event(account_tax_info_id: int, action_name: str): """Emit account payee tax event from passed data. Args: account_tax_info_id (int): Account tax info id. action_name (str): action name. """ if Config.KAFKA_BOOTSTRAP_SERVERS is None: logger.info('Skipping sending Kafka event as no servers specified.') return event_key, event_payload = _create_tax_info_event_key_and_payload( account_tax_info_id ) if action_name not in Config.KAFKA_PRODUCERS_BY_ACTION_NAME: Config.KAFKA_PRODUCERS_BY_ACTION_NAME[action_name] = AccountEventProducer( KAFKA_TOPICS[action_name], ) with Config.KAFKA_PRODUCERS_BY_ACTION_NAME[action_name] as producer: producer.produce_event(event_key, event_payload) def emit_account_change_payee_program_event( account_payee_id: int, payoneer_program_id: int, prev_payoneer_program_id: int ): """Emit account change payee program event from passed data. Args: account_payee_id (int): Account payee id. payoneer_program_id (int): Account new payoneer program id prev_payoneer_program_id (int): Account previous payoneer program id """ if Config.KAFKA_BOOTSTRAP_SERVERS is None: logger.info('Skipping sending Kafka event as no servers specified.') return event_payload = _create_account_payee_payload( account_payee_id, payoneer_program_id, prev_payoneer_program_id ) action_name = ACCOUNT_KAFKA_EVENT_NAMES.CHANGE_PAYEE_PROGRAM if action_name not in Config.KAFKA_PRODUCERS_BY_ACTION_NAME: Config.KAFKA_PRODUCERS_BY_ACTION_NAME[action_name] = AccountEventProducer( KAFKA_TOPICS[action_name], ) with Config.KAFKA_PRODUCERS_BY_ACTION_NAME[action_name] as producer: producer.produce_event(str(account_payee_id), 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_tax_info_event_key_and_payload(account_tax_info_id): """Compose event key and payload.""" key = f'account_tax_info_id_{account_tax_info_id}' payload = { 'accountTaxInfo': { 'accountTaxInfoId': str(account_tax_info_id) }, 'timestamp': datetime.datetime.now().isoformat(), } return key, payload def _create_account_payee_payload( account_payee_id, payoneer_program_id, prev_payoneer_program_id ): """Compose payload.""" payload = { 'accountPayee': { 'accountPayeeId': account_payee_id, 'payoneerProgramId': payoneer_program_id, 'prevPayoneerProgramId': prev_payoneer_program_id } } return payload