"""Unit tests for target event preparation and sending logic.""" from types import SimpleNamespace from uuid import UUID import pytest from src.logic import target_event from src.models.preference_change import SalesforcePreferenceMessage from src.models.preference_change import UpdatedSubscription class TestPrepareMessages: """Tests for prepare_messages function.""" @pytest.fixture def batch_items_sub(self): """Provide batch items with subscription action.""" return { 'sub-123': UpdatedSubscription( id='sub-123', crmId='a0VTy00000H4JxKMAV', mailingListId=UUID('08abc9fb-c4ad-4363-9d7c-8ddc7e3ca778'), mailingListCrmId='a0S61000000ZhGyEAK', emailCampaignId=UUID('edec7008-8147-46ac-b348-8fb430dd3d3e'), oldValue=False, newValue=True, # Subscribe action ) } @pytest.fixture def batch_items_unsub(self): """Provide batch items with unsubscription action.""" return { 'sub-456': UpdatedSubscription( id='sub-456', crmId='a0VTy00000H4JxKMAV', mailingListId=UUID('08abc9fb-c4ad-4363-9d7c-8ddc7e3ca778'), mailingListCrmId='a0S61000000ZhGyEAK', emailCampaignId=UUID('edec7008-8147-46ac-b348-8fb430dd3d3e'), oldValue=True, newValue=False, # Unsubscribe action ) } @pytest.fixture def batch_items_multiple(self): """Provide batch items with multiple subscriptions.""" return { 'sub-123': UpdatedSubscription( id='sub-123', crmId='a0VTy00000H4JxKMAV', mailingListId=UUID('08abc9fb-c4ad-4363-9d7c-8ddc7e3ca778'), mailingListCrmId='a0S61000000ZhGyEAK', emailCampaignId=UUID('edec7008-8147-46ac-b348-8fb430dd3d3e'), oldValue=False, newValue=True, ), 'sub-456': UpdatedSubscription( id='sub-456', crmId='a0VTy00000H4JxLMAV', mailingListId=UUID('18abc9fb-c4ad-4363-9d7c-8ddc7e3ca778'), mailingListCrmId='a0S61000000ZhGzEAK', emailCampaignId=UUID('fdec7008-8147-46ac-b348-8fb430dd3d3e'), oldValue=True, newValue=False, ) } @pytest.fixture def snowflake_response_single(self): """Provide Snowflake response for single subscription.""" return [ { 'FAN_SUBSCRIPTION_ID': 'sub-123', 'SUBSCRIPTION_ID__C': 'sf-sub-123', 'MAILING_LIST_ID__C': 'sf-ml-123', 'FAN_ID__C': 'sf-fan-123', } ] @pytest.fixture def snowflake_response_multiple(self): """Provide Snowflake response for multiple subscriptions.""" return [ { 'FAN_SUBSCRIPTION_ID': 'sub-123', 'SUBSCRIPTION_ID__C': 'sf-sub-123', 'MAILING_LIST_ID__C': 'sf-ml-123', 'FAN_ID__C': 'sf-fan-123', }, { 'FAN_SUBSCRIPTION_ID': 'sub-456', 'SUBSCRIPTION_ID__C': 'sf-sub-456', 'MAILING_LIST_ID__C': 'sf-ml-456', 'FAN_ID__C': 'sf-fan-456', } ] @pytest.fixture def mock_get_subscriptions(self, mocker, snowflake_response_single): """Mock subscription_info.get_subscriptions with default single response.""" return mocker.patch.object( target_event.subscription_info, 'get_subscriptions', return_value=snowflake_response_single ) def test_returns_empty_list_for_empty_batch(self): """Empty batch items should return empty result without calling Snowflake.""" result = target_event.prepare_messages({}) assert result == [] def test_calls_snowflake_with_subscription_ids(self, mock_get_subscriptions, batch_items_sub): """Should extract subscription IDs and fetch data from Snowflake.""" target_event.prepare_messages(batch_items_sub) mock_get_subscriptions.assert_called_once_with(['sub-123']) def test_prepares_salesforce_message_for_subscription(self, mock_get_subscriptions, batch_items_sub): """Should create SalesforcePreferenceMessage with 'Sub' action when new_value is True.""" result = target_event.prepare_messages(batch_items_sub) assert len(result) == 1 message_key, message_value = result[0] assert message_key == 'sf-sub-123' assert isinstance(message_value, SalesforcePreferenceMessage) assert message_value.Subscription_ID__c == 'sf-sub-123' assert message_value.Action__c == 'Sub' assert message_value.Mailing_List_ID__c == 'sf-ml-123' assert message_value.Source__c == 'New Preference Center' assert message_value.Fan_ID__c == 'sf-fan-123' def test_prepares_salesforce_message_for_unsubscription(self, mock_get_subscriptions, batch_items_unsub): """Should create SalesforcePreferenceMessage with 'Unsub' action when new_value is False.""" # Override return value for this specific test mock_get_subscriptions.return_value = [ { 'FAN_SUBSCRIPTION_ID': 'sub-456', 'SUBSCRIPTION_ID__C': 'sf-sub-456', 'MAILING_LIST_ID__C': 'sf-ml-456', 'FAN_ID__C': 'sf-fan-456', } ] result = target_event.prepare_messages(batch_items_unsub) assert len(result) == 1 message_key, message_value = result[0] assert message_key == 'sf-sub-456' assert message_value.Action__c == 'Unsub' def test_prepares_multiple_messages( self, mock_get_subscriptions, batch_items_multiple, snowflake_response_multiple): """Should create multiple SalesforcePreferenceMessage instances for batch items.""" # Override return value for multiple subscriptions mock_get_subscriptions.return_value = snowflake_response_multiple result = target_event.prepare_messages(batch_items_multiple) assert len(result) == 2 # Check first message (subscription) message_key_1, message_value_1 = result[0] assert message_key_1 == 'sf-sub-123' assert message_value_1.Action__c == 'Sub' # Check second message (unsubscription) message_key_2, message_value_2 = result[1] assert message_key_2 == 'sf-sub-456' assert message_value_2.Action__c == 'Unsub' class TestSendMessages: """Tests for send_messages function.""" @pytest.fixture def kafka_config(self, mocker): """Mock Kafka configuration.""" mocker.patch.object( target_event.config, 'settings', SimpleNamespace( kafka_bootstrap_servers='localhost:9092', kafka_security_protocol='PLAINTEXT', kafka_target_topic='test.topic', kafka_message_header_key='test-header' ) ) @pytest.fixture def sample_messages(self): """Provide sample messages to send.""" return [ ( 'sf-sub-123', SalesforcePreferenceMessage( Subscription_ID__c='sf-sub-123', Action__c='Sub', Mailing_List_ID__c='sf-ml-123', Source__c='New Preference Center', Fan_ID__c='sf-fan-123', ) ), ( 'sf-sub-456', SalesforcePreferenceMessage( Subscription_ID__c='sf-sub-456', Action__c='Unsub', Mailing_List_ID__c='sf-ml-456', Source__c='New Preference Center', Fan_ID__c='sf-fan-456', ) ) ] @pytest.fixture def mock_producer(self, mocker, kafka_config): """Create a mock Kafka producer that tracks produced messages.""" class MockProducer: def __init__(self): self.produced_messages = [] def __enter__(self): return self def __exit__(self, exc_type, exc_val, exc_tb): pass def produce(self, topic, event_key, event_value, headers, auto_flush, callback): self.produced_messages.append({ 'topic': topic, 'event_key': event_key, 'event_value': event_value, 'headers': headers, 'auto_flush': auto_flush, 'callback': callback }) mock_producer_instance = MockProducer() mocker.patch.object(target_event, 'EventProducer', return_value=mock_producer_instance) return mock_producer_instance @pytest.fixture def mock_event_producer(self, mocker, kafka_config): """Mock EventProducer class for testing configuration.""" mock_producer_class = mocker.patch.object(target_event, 'EventProducer') mock_producer_class.return_value.__enter__ = mocker.Mock(return_value=mock_producer_class.return_value) mock_producer_class.return_value.__exit__ = mocker.Mock(return_value=None) mock_producer_class.return_value.produce = mocker.Mock() return mock_producer_class def test_creates_event_producer_with_correct_config(self, mock_event_producer, sample_messages): """Should initialize EventProducer with correct configuration.""" target_event.send_messages(sample_messages) mock_event_producer.assert_called_once() call_kwargs = mock_event_producer.call_args[1] assert call_kwargs['bootstrap_servers'] == 'localhost:9092' assert call_kwargs['security_protocol'] == 'PLAINTEXT' assert isinstance(call_kwargs['key_serializer'], target_event.StringSerializer) assert isinstance(call_kwargs['value_serializer'], target_event.SimpleJSONSerializer) def test_produces_messages_to_correct_topic(self, mock_producer, sample_messages): """Should produce all messages to the configured Kafka topic.""" target_event.send_messages(sample_messages) assert len(mock_producer.produced_messages) == 2 for produced_msg in mock_producer.produced_messages: assert produced_msg['topic'] == 'test.topic' def test_produces_messages_with_correct_keys_and_values(self, mock_producer, sample_messages): """Should produce messages with correct keys and serialized values.""" target_event.send_messages(sample_messages) # Check first message first_msg = mock_producer.produced_messages[0] assert first_msg['event_key'] == 'sf-sub-123' assert first_msg['event_value'] == { 'Subscription_ID__c': 'sf-sub-123', 'Action__c': 'Sub', 'Mailing_List_ID__c': 'sf-ml-123', 'Source__c': 'New Preference Center', 'Fan_ID__c': 'sf-fan-123', } # Check second message second_msg = mock_producer.produced_messages[1] assert second_msg['event_key'] == 'sf-sub-456' assert second_msg['event_value'] == { 'Subscription_ID__c': 'sf-sub-456', 'Action__c': 'Unsub', 'Mailing_List_ID__c': 'sf-ml-456', 'Source__c': 'New Preference Center', 'Fan_ID__c': 'sf-fan-456', } def test_produces_messages_with_correct_headers(self, mock_producer, sample_messages): """Should include configured message header in all produced messages.""" target_event.send_messages(sample_messages) for produced_msg in mock_producer.produced_messages: assert produced_msg['headers'] == {'test-header': ''.encode()} def test_produces_messages_without_auto_flush(self, mock_producer, sample_messages): """Should produce messages with auto_flush=False for better performance.""" target_event.send_messages(sample_messages) for produced_msg in mock_producer.produced_messages: assert produced_msg['auto_flush'] is False def test_provides_delivery_callback(self, mock_producer, sample_messages): """Should provide msg_delivery_callback for each produced message.""" target_event.send_messages(sample_messages) for produced_msg in mock_producer.produced_messages: assert produced_msg['callback'] == target_event.msg_delivery_callback class TestMsgDeliveryCallback: """Tests for msg_delivery_callback function.""" @pytest.fixture def mock_message(self): """Create a mock Kafka message object.""" class MockMessage: def key(self): return 'test-key-123' return MockMessage() @pytest.fixture def mock_logger(self, mocker): """Mock the logger for testing log output.""" return mocker.patch.object(target_event, 'logger') def test_raises_error_when_delivery_fails(self, mock_logger, mock_message): """Should log error and re-raise exception when message delivery fails.""" test_error = Exception('Kafka connection failed') with pytest.raises(Exception, match='Kafka connection failed'): target_event.msg_delivery_callback(test_error, mock_message) mock_logger.error.assert_called_once_with( 'There was an error producing message for test-key-123 to Kafka' ) def test_does_not_raise_when_delivery_succeeds(self, mock_logger, mock_message): """Should not raise or log when message is delivered successfully.""" # Should not raise any exception target_event.msg_delivery_callback(None, mock_message) # Should not log any errors mock_logger.error.assert_not_called()