diff --git a/spring-cloud-netflix-turbine-stream/src/main/java/org/springframework/cloud/netflix/turbine/stream/HystrixStreamAggregator.java b/spring-cloud-netflix-turbine-stream/src/main/java/org/springframework/cloud/netflix/turbine/stream/HystrixStreamAggregator.java index 2ab783571..f67bd72b7 100644 --- a/spring-cloud-netflix-turbine-stream/src/main/java/org/springframework/cloud/netflix/turbine/stream/HystrixStreamAggregator.java +++ b/spring-cloud-netflix-turbine-stream/src/main/java/org/springframework/cloud/netflix/turbine/stream/HystrixStreamAggregator.java @@ -17,6 +17,7 @@ package org.springframework.cloud.netflix.turbine.stream; import java.io.IOException; +import java.nio.charset.StandardCharsets; import java.util.List; import java.util.Map; @@ -53,7 +54,8 @@ public class HystrixStreamAggregator { } @ServiceActivator(inputChannel = TurbineStreamClient.INPUT) - public void sendToSubject(@Payload String payload) { + public void sendToSubject(@Payload byte[] bytePayload) { + String payload = new String(bytePayload, StandardCharsets.UTF_8); if (log.isTraceEnabled()) { log.trace("Received hystrix stream payload string: " + payload); } diff --git a/spring-cloud-netflix-turbine-stream/src/test/java/org/springframework/cloud/netflix/turbine/stream/HystrixStreamAggregatorTests.java b/spring-cloud-netflix-turbine-stream/src/test/java/org/springframework/cloud/netflix/turbine/stream/HystrixStreamAggregatorTests.java index 5d9352c39..318c90dec 100644 --- a/spring-cloud-netflix-turbine-stream/src/test/java/org/springframework/cloud/netflix/turbine/stream/HystrixStreamAggregatorTests.java +++ b/spring-cloud-netflix-turbine-stream/src/test/java/org/springframework/cloud/netflix/turbine/stream/HystrixStreamAggregatorTests.java @@ -49,7 +49,7 @@ public class HystrixStreamAggregatorTests { this.publisher.subscribe(map -> { assertThat(map.get("type"), equalTo("HystrixCommand")); }); - this.aggregator.sendToSubject(PAYLOAD); + this.aggregator.sendToSubject(PAYLOAD.getBytes()); this.output.expect(not(containsString("ERROR"))); } @@ -58,7 +58,7 @@ public class HystrixStreamAggregatorTests { this.publisher.subscribe(map -> { assertThat(map.get("type"), equalTo("HystrixCommand")); }); - this.aggregator.sendToSubject("[" + PAYLOAD + "]"); + this.aggregator.sendToSubject(new StringBuilder().append("[").append(PAYLOAD).append("]").toString().getBytes()); this.output.expect(not(containsString("ERROR"))); } @@ -69,7 +69,7 @@ public class HystrixStreamAggregatorTests { }); // If The JSON is embedded in a JSON String this is what it looks like String payload = "\"" + PAYLOAD.replace("\"", "\\\"") + "\""; - this.aggregator.sendToSubject(payload); + this.aggregator.sendToSubject(payload.getBytes()); this.output.expect(not(containsString("ERROR"))); } diff --git a/spring-cloud-netflix-turbine-stream/src/test/java/org/springframework/cloud/netflix/turbine/stream/TurbineStreamTests.java b/spring-cloud-netflix-turbine-stream/src/test/java/org/springframework/cloud/netflix/turbine/stream/TurbineStreamTests.java index db3f02ef9..fb1211cd9 100644 --- a/spring-cloud-netflix-turbine-stream/src/test/java/org/springframework/cloud/netflix/turbine/stream/TurbineStreamTests.java +++ b/spring-cloud-netflix-turbine-stream/src/test/java/org/springframework/cloud/netflix/turbine/stream/TurbineStreamTests.java @@ -22,11 +22,8 @@ import java.io.InputStream; import java.io.InputStreamReader; import java.net.URI; import java.util.Map; - -import com.fasterxml.jackson.databind.ObjectMapper; import org.junit.Test; import org.junit.runner.RunWith; - import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; @@ -34,6 +31,10 @@ import org.springframework.boot.test.context.SpringBootTest; import org.springframework.boot.web.server.LocalServerPort; import org.springframework.cloud.contract.stubrunner.StubTrigger; import org.springframework.cloud.contract.stubrunner.spring.AutoConfigureStubRunner; +import org.springframework.cloud.contract.verifier.messaging.MessageVerifier; +import org.springframework.cloud.contract.verifier.messaging.stream.StreamStubMessages; +import org.springframework.context.ApplicationContext; +import org.springframework.context.annotation.Bean; import org.springframework.http.HttpHeaders; import org.springframework.http.HttpMethod; import org.springframework.http.HttpRequest; @@ -44,9 +45,11 @@ import org.springframework.http.client.ClientHttpRequestExecution; import org.springframework.http.client.ClientHttpRequestInterceptor; import org.springframework.http.client.ClientHttpResponse; import org.springframework.integration.support.management.MessageChannelMetrics; +import org.springframework.messaging.Message; import org.springframework.messaging.SubscribableChannel; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; import org.springframework.web.client.RestTemplate; +import com.fasterxml.jackson.databind.ObjectMapper; import static org.assertj.core.api.Assertions.assertThat; import static org.springframework.boot.test.context.SpringBootTest.WebEnvironment.RANDOM_PORT; @@ -61,7 +64,7 @@ import static org.springframework.boot.test.context.SpringBootTest.WebEnvironmen // https://github.com/spring-cloud/spring-cloud-netflix/issues/1948 "spring.cloud.stream.bindings.turbineStreamInput.destination=hystrixStreamOutput", "spring.jmx.enabled=true", "stubrunner.workOffline=true", - "stubrunner.ids=org.springframework.cloud:spring-cloud-netflix-hystrix-stream:${projectVersion:2.0.0.BUILD-SNAPSHOT}:stubs" }) + "stubrunner.ids=org.springframework.cloud:spring-cloud-netflix-hystrix-stream:${projectVersion:2.0.0.BUILD-SNAPSHOT}:stubs"}) @AutoConfigureStubRunner public class TurbineStreamTests { @Autowired @@ -85,6 +88,22 @@ public class TurbineStreamTests { @EnableAutoConfiguration @EnableTurbineStream public static class Application { + @Bean + //TODO This can be removed after Finchley.RELEASE, once we can use Spring Cloud Contract Verifier 2.0.0 + //This is a hack to allow compatibility between Stream 2.0.0, which is sending everything as a byte array, + //and contract which is assuming everything is a String. + public MessageVerifier> customMessageVerifier(ApplicationContext context) { + return new StreamStubMessages(context) { + @Override + public void send(T payload, Map headers, String destination) { + if(String.class.isInstance(payload)){ + super.send(((String)payload).getBytes(), headers, destination); + return; + } + super.send(payload, headers, destination); + } + }; + } } @Test