"""Interface for work with events.""" import logging from time import time from confluent_kafka.serialization import StringSerializer from ddtrace import tracer from flask import g from flask import has_app_context from kafka_utils.producer.event import EventProducer from kafka_utils.producer.serializer.avro import build_avro_serializer from product_digital import config logging.basicConfig(level=logging.INFO) logger = logging.getLogger(__name__) class OperationType: """Enum for operations with products.""" CREATE = 'create' UPDATE = 'update' class OperationContext: """Enum for operation context.""" NEW = 'new' COPY = 'copy' SUBMIT = 'submit' UNSUBMIT = 'unsubmit' APPROVE = 'approve' string_serializer = StringSerializer() if config.ENVIRONMENT != config.TEST_ENVIRONMENT: if not config.SCHEMA_REGISTRY_URL: raise ValueError('"SCHEMA_REGISTRY_URL" variable is not configured!') if not config.KAFKA_DIGITAL_PRODUCT_TOPIC: raise ValueError('"KAFKA_DIGITAL_PRODUCT_TOPIC" variable is not configured!') if not config.DIGITAL_PRODUCT_SCHEMA_VERSION: raise ValueError('"DIGITAL_PRODUCT_SCHEMA_VERSION" variable is not configured!') logger.info('Creating digital product event avro serializer...') avro_serializer = build_avro_serializer( config.SCHEMA_REGISTRY_URL, config.KAFKA_DIGITAL_PRODUCT_TOPIC, config.DIGITAL_PRODUCT_SCHEMA_VERSION) logger.info('The digital product event avro serializers is created') logger.info('Creating digital product event producer...') digital_product_producer = EventProducer( config.KAFKA_BOOTSTRAP_SERVERS, string_serializer, avro_serializer, 'SSL') logger.info('digital product event producer is created') else: digital_product_producer = None _OPTIONAL_PAYLOAD_PROPERTIES = ('upc', 'product_name', 'project_id') def _create_product_event_key_and_payload(product_id, **payload_properties): """Compose event key and payload.""" key = f'product_id_{product_id}' payload = dict(product_id=product_id) for payload_property in _OPTIONAL_PAYLOAD_PROPERTIES: if payload_property in payload_properties: payload[payload_property] = payload_properties[payload_property] return key, payload @tracer.wrap() def emit_digital_product_event( *, product_id, operation_type, operation_context=None, **payload_properties): """Emit product event from passed data. Args: product_id (int): Product unique identifier. operation_type (str): Enum from OperationType. operation_context (str): Enum from OperationContext. **payload_properties: Optional product properties from _OPTIONAL_PAYLOAD_PROPERTIES """ event_key, event_payload = _create_product_event_key_and_payload( product_id, **payload_properties) identity_id = None if has_app_context() and hasattr(g, 'request_context') and hasattr(g.request_context, 'identity_uuid'): identity_id = g.request_context.identity_uuid event_value = { 'operation': { 'type': operation_type, 'timestamp': time() * 1000, 'context': operation_context, 'identity_id': identity_id }, 'payload': event_payload } with digital_product_producer: digital_product_producer.produce( config.KAFKA_DIGITAL_PRODUCT_TOPIC, event_key, event_value, callback=_default_on_delivery) logger.info(f'Successfully produced event with {event_key} key') def _default_on_delivery(err, msg): if err is not None: g.ows.log.error( 'Delivery failed for Product record {}: {}'.format(msg.key(), err)) elif msg: g.ows.log.info( 'Product record {} successfully produced to {} [{}] at offset {}'.format( msg.key(), msg.topic(), msg.partition(), msg.offset()))