"""Kafka test module.""" from src.common.exceptions import exceptions from src.common.connectors import kafka_producer import json from unittest.mock import MagicMock from unittest.mock import patch import pytest @patch('src.common.connectors.kafka_producer.logger') def test_delivery_report(mock_logger): """Test delivery report method.""" msg = MagicMock() msg.topic.return_value = 'fake-topic' msg.partition.return_value = '0' kafka_producer.delivery_report(None, msg) mock_logger.info.asset_called_with( 'Message delivered to fake-topic [0]\n') @patch('src.common.connectors.kafka_producer.logger') def test_delivery_report_exception(mock_logger): """Test delivery report method exception.""" msg = MagicMock() msg.topic.return_value = 'fake-topic' msg.partition.return_value = '0' with pytest.raises(exceptions.RetryableException): kafka_producer.delivery_report(KeyError, msg) mock_logger.exception.assert_called_with( "Message delivery failed: ") @patch('src.common.connectors.kafka_producer.delivery_report') @patch('src.common.connectors.kafka_producer.p') def test_produce_message( mock_producer, mock_delivery_report, mock_delivery_event): """Test produce_message method.""" mock_data = json.dumps(mock_delivery_event) kafka_producer.produce_message(mock_data, 'fake-topic') assert mock_producer.produce.called assert mock_producer.flush.called mock_producer.produce.assert_called_with( 'fake-topic', mock_data.encode('utf-8'), callback=mock_delivery_report)