"""The main Kafka message processing logic.""" import sentry_sdk import simplejson as json from neo4j_sync_dlq.connectors import logger def process(message): """Process the message - log and optionally send back to Kafka. Args: message (KafkaMessage): parsed kafka message. """ extra = message.to_sentry_message() # TODO: implement the retry logic log = logger.get_current_logger() log.info(json.dumps(extra)) with sentry_sdk.push_scope() as scope: for key, value in extra.items(): scope.set_extra(key, value) sentry_sdk.capture_message( f'Neo4J Sync DLQ message for {message.table} table')