"""Test JSON Schema messages.""" from kafka_utils.consumer.deserializer.string import StringDeserializer from kafka_utils.consumer.source.mapping import EventSourceMessage def test_consume_json_schema_message(mock_json_schema_event, mock_json_schema_serializer): """Test decoding a JSON schema message.""" string_deserializer = StringDeserializer() event = EventSourceMessage(mock_json_schema_event) for batch_key, kafka_event in event: assert batch_key == 'event.owsContentReview2.review-2' assert string_deserializer.deserialize(kafka_event.key) == 'example-product-1' mock_event_value = mock_json_schema_serializer.deserialize( kafka_event.value, kafka_event.topic, 'value') assert mock_event_value == { 'operation': { 'context': 'new', 'timestamp': 1655390557, 'type': 'create' }, 'payload': { 'product_id': 1, 'review_note': '', 'review_queue_id': 1, 'user_id': 'linus' }} mock_event_value_schema_id = mock_json_schema_serializer.get_schema_id(kafka_event.value) assert mock_event_value_schema_id == 81