"""Test debezium message.""" import json import pytest from kafka_utils.consumer.message.debezium import DebeziumMessage from kafka_utils.consumer.source.mapping import EventSourceMessage from kafka_utils.exceptions import DebeziumMessageException from kafka_utils.exceptions import IneligibleEventException def test_debezium_message(): """Test DebeziumMessage.""" assert DebeziumMessage( json.dumps({'payload': {'op': 'c', 'before': None, 'after': {}}}), 'whatever') assert DebeziumMessage( json.dumps({'payload': {'op': 'd', 'before': {}, 'after': None}}), 'whatever') def test_debezium_message_get_field(mock_event): """Test DebeziumMessage.""" mock_event_source_message = EventSourceMessage( mock_event) events = [x for _, x in mock_event_source_message] # Test create event create_dm = DebeziumMessage( events[0].value, events[0].topic) assert create_dm.operation == 'c' assert create_dm.after assert create_dm.before is None field_transition = create_dm.get_field_values('release_status') assert field_transition == { 'before': None, 'after': 'orchard_processing' } # Test update event update_dm = DebeziumMessage( events[1].value, events[1].topic) assert update_dm.operation == 'u' field_transition = update_dm.get_field_values('release_status') assert field_transition == { 'before': 'orchard_processing', 'after': 'orchard_processing' } def test_debezium_message_raises(): """Test DebeziumMessage raises Exception.""" with pytest.raises(IneligibleEventException): DebeziumMessage( message=json.dumps({'payload': {'op': 'u'}}), allowed_event_ops=['c'], topic='whatever') with pytest.raises(IneligibleEventException): DebeziumMessage( message=None, topic='whatever') with pytest.raises(DebeziumMessageException): DebeziumMessage( message=json.dumps({}), topic='whatever') with pytest.raises(DebeziumMessageException): DebeziumMessage( message=json.dumps({'payload': None}), topic='whatever') with pytest.raises(DebeziumMessageException): DebeziumMessage( message=json.dumps({'payload': {'op': None}}), topic='whatever') with pytest.raises(DebeziumMessageException): DebeziumMessage( message=json.dumps({'payload': {'op': 'u', 'before': None, 'after': {}}}), topic='whatever') with pytest.raises(DebeziumMessageException): DebeziumMessage( json.dumps({'payload': {'op': 'u', 'before': {}, 'after': None}}), topic='whatever') def test_debezium_events_from_msk_message( mock_multi_topic_msk_message): """Test DebeziumEventMessage.""" message = EventSourceMessage(mock_multi_topic_msk_message) for key, event in message: if event.topic.startswith('cdc.'): review_event = DebeziumMessage( message=event.value, topic=event.topic) assert review_event.message['op'] assert review_event.message['payload']