"""Product event logging adapter.""" from kafka_utils.consumer.message.product import ProductEventMessage from kafka_utils.consumer.source.msk import MSKMessage from kafka_utils.exceptions import EventProducerMessageException from content_utils.logging.logger import ContentLambdaLogger class ProductEventLogger(ContentLambdaLogger): """Logger for product events.""" def start(self, msk_message: MSKMessage = None, event_key=None, deserializer=None): """Begin processing.""" super().start(msk_message, event_key) if deserializer is None: return try: event = ProductEventMessage(msk_message.value, msk_message.topic, deserializer) except EventProducerMessageException: return op_pieces = [] if event.operation_type: op_pieces.append(event.operation_type) if event.operation_context: op_pieces.append(event.operation_context) self.data.update( operation=':'.join(op_pieces) if len(op_pieces) else '', product_id=event.product_id or '', queue_id=event.review_queue_id or '' )