"""Processing logic.""" from content_utils.exceptions import InvalidMessageException from content_utils.exceptions import InvalidProductException from content_utils.exceptions import NoProductDataException from kafka_utils.consumer.message.product import ProductEventMessage from kafka_utils.exceptions import IneligibleEventException from src import config from src import constants from src.logic.add_product_to_queue import add_product_to_queue from src.logic.extract_message_payload import extract_product_message_body from src.logic.filter_ineligible_products import filter_ineligible_products SKIPPED_EXCEPTIONS = ( InvalidMessageException, InvalidProductException, IneligibleEventException, NoProductDataException ) def handle_product_submission(product_msg, product_id, correlation_id): """Submit a product into the review_queue db.""" if product_msg.operation_context != 'submit': raise IneligibleEventException(f'skip operations {product_msg.operation_context}') if product_id in config.EXPLICIT_SKIP_PRODUCT_IDS: raise IneligibleEventException('skip product id based on configuration') product = extract_product_message_body(product_id, correlation_id) filter_ineligible_products(product) identity_id = product_msg.message.get('operation', {}).get('identity_id') submission_datetime = None timestamp_str = product_msg.message.get('operation', {}).get('timestamp') if timestamp_str: submission_datetime = timestamp_str.strftime('%Y-%m-%d %H:%M:%S') review_queue_id = add_product_to_queue( product_id, correlation_id, identity_id, submission_datetime) return review_queue_id def process_event(topic, value, deserializer, correlation_id): """Single event handling logic.""" product_id = None try: product_msg = ProductEventMessage(value, topic, deserializer) product_id = product_msg.product_id review_queue_id = handle_product_submission(product_msg, product_id, correlation_id) output_msg = f'Product {product_id} added to review queue {review_queue_id}' return { 'status': constants.STATUS_OK, 'product_id': product_id, 'review_queue_id': review_queue_id, 'message': output_msg, } except SKIPPED_EXCEPTIONS as e: return { 'status': constants.STATUS_SKIP, 'product_id': product_id, 'message': str(e), } except Exception as e: return { 'status': constants.STATUS_ERROR, 'product_id': product_id, 'message': str(e) }