"""Kafka consumer tests.""" import json from unittest.mock import patch from unittest.mock import PropertyMock from neo4j_sync_dlq.logic import kafka from neo4j_sync_dlq.logic import kafka_message @patch.object(kafka.processing, 'process', wraps=kafka.processing.process) @patch.object(kafka.processing.sentry_sdk, 'capture_message') def test_kafka_consumer_send_message_to_sentry( sentry_spy, processing_spy, mock_kafka_consumer, mock_kafka_message): """Test kafka consumer receive the message and send to sentry.""" with patch( 'neo4j_sync_dlq.logic.kafka.KafkaWorker', wraps=kafka.KafkaWorker) as KafkaWorkerPatch: worker = KafkaWorkerPatch('dlq_topic', 'group_id', {}) do_work_mock = PropertyMock(side_effect=[True, False]) type(worker).do_work = do_work_mock worker.run() processing_spy.assert_called() assert processing_spy.call_args[0][0].message == json.loads( mock_kafka_message.value.decode()) sentry_spy.assert_called_with( 'Neo4J Sync DLQ message for test_table table') def test_kafka_message(mock_kafka_message): """Test kafka message instance is created from tge received data.""" message = kafka_message.KafkaMessage(mock_kafka_message) assert message.database == 'test_db' assert message.table == 'test_table' assert message.data == { 'artist_id': 123, 'artist_name': 'test_artist', 'vendor_id': 234 } assert message.error_exception_name == 'TestExceptionClassName' assert message.error_exception_message == 'test_exception_message'