"""Test event source mapping.""" import pytest from kafka_utils.consumer.deserializer.simple_json import JSONDeserializer from kafka_utils.consumer.source.mapping import EventSourceMessage from kafka_utils.exceptions import MSKMessageException def test_event_source_message(mock_event): """Test LambdaMSKEventSourceMessage.""" json_deserializer = JSONDeserializer() mock_event_source_message = EventSourceMessage(mock_event) for key, event in mock_event_source_message: deserialized_event = json_deserializer.deserialize(event.value) assert 'cdc.artRelations.releases-' in key assert 'payload' in deserialized_event assert deserialized_event.get('schema').get('type') == 'struct' def test_multi_topic_event_source_message(mock_multi_topic_msk_message): """Test LambdaMSKEventSourceMessage with multiple topics.""" mock_message = EventSourceMessage( mock_multi_topic_msk_message) for _, event in mock_message: assert event.topic.startswith( 'cdc') or event.topic.startswith('event') assert event.value def test_invalid_event_source_message(): """Test LambdaMSKEventSourceMessage raises with an invalid event.""" with pytest.raises(MSKMessageException): EventSourceMessage({}) def test_empty_event_source_message(mock_empty_event): """Test generator returns None for empty messages.""" mock_event_source_message = EventSourceMessage(mock_empty_event) for key, event in mock_event_source_message: assert 'cdc.artRelations.releases-' in key assert event.value is None