"""Interface for work with events.""" import json import logging from confluent_kafka import Producer, error as kafka_error from abacus_contract.connectors import sentry from abacus_contract.constants.constants import CONTRACT_KAFKA_EVENT_TYPES, KAFKA_TOPICS from abacus_contract.utils.exception import KafkaSendEventException from core.config import Config logger = logging.getLogger('kafka_producer_logger') class ContractProducer: """Contract event producer.""" def __init__(self, topic): """Initialize ContractProducer.""" logger.info('Initializing ContractProducer') 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('ContractProducer 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: logger.debug('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, ) logger.debug('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_contract_event(contract_id: int, action_name: str): """Emit contract event from passed data. Args: contract_id (int): Contract id. action_name (str): action name. """ if not Config.KAFKA_BOOTSTRAP_SERVERS: logger.info('Skipping sending Kafka event as no servers specified.') return event_key, event_payload = _create_contract_event_key_and_payload( contract_id, action_name ) if action_name not in Config.KAFKA_PRODUCERS_BY_ACTION_NAME: Config.KAFKA_PRODUCERS_BY_ACTION_NAME[action_name] = ContractProducer( 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 event `{msg.key()}`: {str(err)}', err, 'error', f'Error! Code: {type(err).__name__}, Message, {str(err)}', ) raise KafkaSendEventException( f'Delivery failed for event `{msg.key()}`: {str(err)}' ) elif msg: sentry.send_to_sentry( 'Event record {} successfully produced to {} [{}] at offset {}'.format( msg.key(), msg.topic(), msg.partition(), msg.offset() ), {}, 'info', msg, ) def _create_contract_event_key_and_payload(contract_id: int, action_name: str): """Compose event key and payload.""" key = f'contract_id_{contract_id}' payload = { 'event_type': CONTRACT_KAFKA_EVENT_TYPES[action_name], 'contract_id': contract_id, } return key, payload