"""CDC event logging adapter.""" from kafka_utils.consumer.message.debezium import DebeziumMessage from kafka_utils.exceptions import DebeziumMessageException from content_utils.logging.logger import ContentLambdaLogger class ContentLambdaLoggerCDC(ContentLambdaLogger): """Logger for `cdc.*` (mysql/debezium) events.""" def start(self, msk_message=None, event_key=None): """Begin processing.""" super().start(msk_message, event_key) if msk_message is None or msk_message.value is None: return try: db_msg = DebeziumMessage( message=msk_message.value, topic=msk_message.topic) except DebeziumMessageException: return product_id = self.get_product_id_from_event(db_msg) self.data.update(product_id=product_id, operation=db_msg.operation) queue_id = self.get_queue_id_from_event(db_msg, event_key) if queue_id: self.data.update(queue_id=queue_id) def get_product_id_from_event(self, db_msg): """Get product_id out of debezium message.""" if db_msg.after: if 'release_id' in db_msg.after: return db_msg.after['release_id'] if 'product_id' in db_msg.after: return db_msg.after['product_id'] if db_msg.before: if 'release_id' in db_msg.before: return db_msg.before['release_id'] if 'product_id' in db_msg.before: return db_msg.before['product_id'] def get_queue_id_from_event(self, db_msg, event_key): """Get queue_id out of debezium message.""" if not event_key.startswith('cdc.contentReview.reviewQueue-'): return if db_msg.after and 'id' in db_msg.after: return db_msg.after['id'] if db_msg.before and 'id' in db_msg.before: return db_msg.before['id']