diff --git a/spring-cloud-contract-verifier/src/main/java/org/springframework/cloud/contract/verifier/messaging/kafka/ContractVerifierKafkaConfiguration.java b/spring-cloud-contract-verifier/src/main/java/org/springframework/cloud/contract/verifier/messaging/kafka/ContractVerifierKafkaConfiguration.java index 959533fd80..ed9fab3055 100644 --- a/spring-cloud-contract-verifier/src/main/java/org/springframework/cloud/contract/verifier/messaging/kafka/ContractVerifierKafkaConfiguration.java +++ b/spring-cloud-contract-verifier/src/main/java/org/springframework/cloud/contract/verifier/messaging/kafka/ContractVerifierKafkaConfiguration.java @@ -16,6 +16,10 @@ package org.springframework.cloud.contract.verifier.messaging.kafka; +import java.nio.charset.StandardCharsets; +import java.util.Map; +import java.util.stream.Collectors; + import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; @@ -35,6 +39,7 @@ import org.springframework.context.annotation.Configuration; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.test.EmbeddedKafkaBroker; import org.springframework.messaging.Message; +import org.springframework.messaging.MessageHeaders; /** * @author Marcin Grzejszczak @@ -80,7 +85,19 @@ class ContractVerifierKafkaHelper extends ContractVerifierMessaging> @Override protected ContractVerifierMessage convert(Message message) { - return new ContractVerifierMessage(message.getPayload(), message.getHeaders()); + return new ContractVerifierMessage(message.getPayload(), convertHeaders(message.getHeaders())); + } + + private MessageHeaders convertHeaders(Map headers) { + return new MessageHeaders(headers.entrySet().stream() + .collect(Collectors.toMap(Map.Entry::getKey, e -> maybeConvertValue(e.getValue())))); + } + + private Object maybeConvertValue(Object value) { + if (!(value instanceof byte[])) { + return value; + } + return new String((byte[]) value, StandardCharsets.UTF_8); } } diff --git a/spring-cloud-contract-verifier/src/test/groovy/org/springframework/cloud/contract/verifier/messaging/kafka/ContractVerifierKafkaHelperSpec.groovy b/spring-cloud-contract-verifier/src/test/groovy/org/springframework/cloud/contract/verifier/messaging/kafka/ContractVerifierKafkaHelperSpec.groovy new file mode 100644 index 0000000000..550f2fa51f --- /dev/null +++ b/spring-cloud-contract-verifier/src/test/groovy/org/springframework/cloud/contract/verifier/messaging/kafka/ContractVerifierKafkaHelperSpec.groovy @@ -0,0 +1,77 @@ +/* + * Copyright 2013-2020 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cloud.contract.verifier.messaging.kafka + +import spock.lang.Specification + +import org.springframework.cloud.contract.verifier.messaging.MessageVerifier +import org.springframework.cloud.contract.verifier.messaging.internal.ContractVerifierMessage +import org.springframework.cloud.contract.verifier.messaging.kafka.ContractVerifierKafkaHelper +import org.springframework.messaging.Message +import org.springframework.messaging.support.MessageBuilder + +import static org.mockito.Mockito.mock + +/** + * Unit tests for {@link ContractVerifierKafkaHelper}. + * + * @author Chris Bono + */ +class ContractVerifierKafkaHelperSpec extends Specification { + + def "should convert message with no headers"() { + given: + Message message = MessageBuilder + .withPayload("some-data") + .build() + ContractVerifierKafkaHelper contractVerifierKafkaHelper = new ContractVerifierKafkaHelper(mock(MessageVerifier.class)) + when: + ContractVerifierMessage contractVerifierMessage = contractVerifierKafkaHelper.convert(message) + then: + contractVerifierMessage.payload == "some-data" + } + + def "should convert message with basic header"() { + given: + Message message = MessageBuilder + .withPayload("some-data") + .setHeader("some-header", "5150") + .build() + ContractVerifierKafkaHelper contractVerifierKafkaHelper = new ContractVerifierKafkaHelper(mock(MessageVerifier.class)) + when: + ContractVerifierMessage contractVerifierMessage = contractVerifierKafkaHelper.convert(message) + then: + contractVerifierMessage.payload == "some-data" + contractVerifierMessage.headers.containsKey("some-header") + contractVerifierMessage.headers.get("some-header") == "5150" + } + + def "should convert message with byte[] header"() { + given: + Message message = MessageBuilder + .withPayload("some-data") + .setHeader("some-header", "5150".getBytes()) + .build() + ContractVerifierKafkaHelper contractVerifierKafkaHelper = new ContractVerifierKafkaHelper(mock(MessageVerifier.class)) + when: + ContractVerifierMessage contractVerifierMessage = contractVerifierKafkaHelper.convert(message) + then: + contractVerifierMessage.payload == "some-data" + contractVerifierMessage.headers.containsKey("some-header") + new String(contractVerifierMessage.headers.get("some-header")) == "5150" + } +}