"""Module that responds for input message processing.""" from enum import Enum from ddtrace import tracer from switchboard_consumer import kafka_message from switchboard_consumer.connectors import log_central from switchboard_consumer.connectors.sentry import sentry_client from switchboard_consumer.constants import entity_type, message_type, system from switchboard_consumer.formatters.errors import format_exception from switchboard_consumer.logic.product_metadata_update import \ handle_product_metadata_update from switchboard_consumer.logic.project_metadata_update import \ handle_project_metadata_update from switchboard_consumer.logic.track_metadata_update import \ handle_track_metadata_update from switchboard_consumer.utils.datadog import setup_datadog_tags from switchboard_consumer.utils.message import create_processing_result handlers = { entity_type.PROJECT: { message_type.METADATA_UPDATE: handle_project_metadata_update, }, entity_type.PRODUCT: { message_type.METADATA_UPDATE: handle_product_metadata_update }, entity_type.TRACK: { message_type.METADATA_UPDATE: handle_track_metadata_update } } class ReturnCode(Enum): """Return codes for processing messages.""" BAD_MESSAGE_FORMAT = 1 INCORRECT_SENDING_SYSTEM = 2 UNKNOWN_MESSAGE_TYPE = 3 def __bool__(self): """Return False for return codes.""" return False @tracer.wrap(service='daemon-switchboard-consumer') def process(message, switchboard_client, orchard_client): """ Process a message. :param message: JSON :param switchboard_client: :param orchard_client: :return: """ logger = log_central.get_current_logger() try: message = kafka_message.KafkaMessage(message) logger = log_central.get_current_logger(message.correlation_id) logger.info( 'Processing new message') except (kafka_message.BadMessageFormat, AttributeError): logger.warning('Bad message format. Message ignored.') return ReturnCode.BAD_MESSAGE_FORMAT if message.sending_system != system.SONY: logger.info('Ignoring own message') return ReturnCode.INCORRECT_SENDING_SYSTEM logger.info('Received message {}:{}'.format( message.entity_type, message.message_type), custom_fields={'raw_message': message.message}) setup_datadog_tags(tracer, message, logger) try: handler = handlers[message.entity_type][message.message_type] logger.info('Identified message handler for: {}:{}' .format(message.entity_type, message.message_type)) except KeyError: logger.warning('Unknown message type. Message ignored.') return ReturnCode.UNKNOWN_MESSAGE_TYPE try: responses = handler( logger, message, switchboard_client, orchard_client) logger.info('Successfully processed message for: {}:{}'.format( message.entity_type, message.message_type)) for response in responses: logger.info('Returning message: {}:{}' .format(response['entityType'], response['messageType']), custom_fields={'raw_message': response}) return responses except Exception as err: if sentry_client: sentry_client.captureException() logger.error( 'Failed to process message with error: {error}'.format( error=repr(err))) logger.exception(err) return [create_processing_result(message.ids, message, exception=format_exception(err))]