"""Unit tests for the kafka_producer module.""" from confluent_kafka import KafkaError from confluent_kafka import Message import pytest import config # noqa from src import kafka_producer from src.models import TriggeredSendFailedEvent from src.models import TriggeredSendSuccessfulEvent @pytest.fixture def failed_event(): """Sample TriggeredSendFailedEvent object for testing.""" return TriggeredSendFailedEvent( triggered_send_info={ 'business_unit_id': 'test-business-unit-id', 'triggered_send_definition_key': 'test-definition-key', 'recipients': [ { 'email_address': 'test@example.com', 'profile_uuid': 'test-profile-uuid', 'source_event_id': 'test-event-id', 'first_name': 'Test', 'subscriber_key': 'test-subscriber-key', 'last_name': 'User', 'postal_code': '12345' } ] }, error='Test error message', api_response=None ) @pytest.fixture def successful_event(): """Sample TriggeredSendSuccessfulEvent object for testing.""" return TriggeredSendSuccessfulEvent( triggered_send_info={ 'business_unit_id': 'test-business-unit-id', 'triggered_send_definition_key': 'test-definition-key', 'recipients': [ { 'email_address': 'test@example.com', 'profile_uuid': 'test-profile-uuid', 'source_event_id': 'test-event-id', 'first_name': 'Test', 'subscriber_key': 'test-subscriber-key', 'last_name': 'User', 'postal_code': '12345' } ] } ) class TestTriggeredSendEventProducer: """Tests for the TriggeredSendEventProducer class.""" def test_produce_successful_event(self, mocker, successful_event): """Test producing a successful event.""" # Setup triggered_send_event_producer = kafka_producer.TriggeredSendEventProducer() triggered_send_event_producer.producer = mocker.Mock() # Execute triggered_send_event_producer.produce_successful_event(successful_event) # Verify triggered_send_event_producer.producer.produce.assert_called_once_with( topic=config.SUCCESSFUL_TRIGGERED_SENDS_TOPIC, event_key=None, event_value=successful_event.model_dump(), callback=kafka_producer._msg_delivery_callback ) def test_produce_failed_event(self, mocker, failed_event): """Test producing a failed event.""" # Setup triggered_send_event_producer = kafka_producer.TriggeredSendEventProducer() triggered_send_event_producer.producer = mocker.Mock() # Execute triggered_send_event_producer.produce_failed_event(failed_event) # Verify triggered_send_event_producer.producer.produce.assert_called_once_with( topic=config.FAILED_TRIGGERED_SENDS_TOPIC, event_key=None, event_value=failed_event.model_dump(), callback=kafka_producer._msg_delivery_callback ) class TestMsgDeliveryCallback: """Tests for the _msg_delivery_callback function.""" def test_msg_delivery_raises_exception(self, mocker): """Test that an exception is raised when there is an error.""" # Setup mock_message = mocker.Mock(spec=Message) mock_message.topic.return_value = 'test-topic' mock_error = mocker.Mock(spec=KafkaError) mock_error.__str__ = mocker.Mock(return_value='Test Kafka error') # Execute and Verify with pytest.raises(Exception, match='Kafka producer error: Test Kafka error'): kafka_producer._msg_delivery_callback(mock_error, mock_message) def test_msg_delivery_logs_debug_message(self, mocker): """Test that a debug message is logged when there is no error.""" # Setup mock_message = mocker.Mock(spec=Message) mock_message.topic.return_value = 'test-topic' mock_message.partition.return_value = 0 mock_message.offset.return_value = 123 mock_logger = mocker.patch('src.kafka_producer.logger') # Execute kafka_producer._msg_delivery_callback(None, mock_message) # Verify mock_logger.debug.assert_called_once_with( 'Message delivered to test-topic [0] at offset 123' )