"""Interface for work with events.""" import json import logging from confluent_kafka import error as kafka_error from confluent_kafka import Producer from abacus_contract.config import Config 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 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 Config.KAFKA_BOOTSTRAP_SERVERS is None: 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