"""Unit tests for dlq_event helpers.""" from unittest.mock import call from src.logic import dlq_event class TestSendInvalidMessages: """Tests for send_invalid_messages.""" def test_returns_early_when_messages_absent(self, dlq_settings, event_producer_mocks): """Should not invoke EventProducer when no invalid messages exist.""" dlq_event.send_invalid_messages([]) event_producer_mocks.cls.assert_not_called() def test_produces_messages_with_expected_headers(self, dlq_settings, sample_invalid_messages, event_producer_mocks): """Should publish each invalid message with encoded error context.""" dlq_event.send_invalid_messages(sample_invalid_messages) assert event_producer_mocks.string_serializer.call_count == 2 event_producer_mocks.cls.assert_called_once_with( bootstrap_servers=dlq_settings.kafka_bootstrap_servers, key_serializer=event_producer_mocks.key_serializer, value_serializer=event_producer_mocks.value_serializer, security_protocol=dlq_settings.kafka_security_protocol, ) assert event_producer_mocks.context_manager.__enter__.called expected_calls = [ call( topic=dlq_settings.kafka_dlq_topic, event_key=sample_invalid_messages[0].message_key, event_value=sample_invalid_messages[0].message_value, headers={ 'errorType': sample_invalid_messages[0].error_type.encode(), 'errorMessage': sample_invalid_messages[0].error_message.encode(), }, auto_flush=False, ), call( topic=dlq_settings.kafka_dlq_topic, event_key=sample_invalid_messages[1].message_key, event_value=sample_invalid_messages[1].message_value, headers={ 'errorType': sample_invalid_messages[1].error_type.encode(), 'errorMessage': sample_invalid_messages[1].error_message.encode(), }, auto_flush=False, ), ] event_producer_mocks.producer.produce.assert_has_calls(expected_calls)