Content type redesign

Fixes #992, #1050, #1051, #1052

Adding custom jackson converter with some tests
Adds kryo message converter to replace codec
Checkstyle changes

Removing codec support
- Removed codec dependency from AbstractBinder
- MessageSerializationUtils is almost an empty shell for now, just to
  keep code compiling until we get EmbeddedHeaders interceptors
- Updated Kryo tests

Removing codec module from build
Added a new Annotation for custom converters '@StreamConverter'
Fixed some tests with new expected behavior
Moved broken tests to a temporary package to keep track of progress
Fixed KryoConverter to fail based on headers
Fixed a couple of more tests

Making converters strict to only convert their corresponding contentType
Bypassing conversion for ErrorMessages

* Configuring SI ConfigurableCompositeMessageConverter
 - Moved ContentType related beans into separate configuration
 - Configured SI ConfigurableCompositeMessageConverter to use same
   converters as Stream does (for ServiceActivator)
- TupleConverter should return byte[] as all other converters
- Fixed tests

* Fixes tests
 - Revert to Boot 2.0.0.M3. Snapshots breaking actuator
 - Checkstyle fixes
 - Disable JsonUnmarshalling as a catch all converter

Fixing Schema tests
Fixing Metrics tests
Fixing reactive tests
applying checkstyle fixes

 * Adding new content type tests
 - Fixed ContentTypeInterceptor misusage of default mimeType

Changing contentType doc section
Improving doc section
Last minute polish
Fixing BinderTests to use bytes to compare messages
Applied changes to Base Binders test to use the new contentType handling mechanism
PR review fixes

Renaming StreamConverter -> StreamMessageConverter
This commit is contained in:
Vinicius Carvalho
2017-09-05 13:45:21 -04:00
committed by Soby Chacko
parent 24cf992301
commit 171f034a8c
74 changed files with 1745 additions and 1066 deletions

View File

@@ -27,6 +27,7 @@ import org.junit.Test;
import org.reactivestreams.Publisher;
import reactor.core.publisher.Flux;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
import org.springframework.cloud.stream.annotation.EnableBinding;
import org.springframework.cloud.stream.annotation.Output;
@@ -34,7 +35,6 @@ import org.springframework.cloud.stream.messaging.Processor;
import org.springframework.cloud.stream.messaging.Source;
import org.springframework.cloud.stream.test.binder.MessageCollector;
import org.springframework.context.ConfigurableApplicationContext;
import org.springframework.context.annotation.AnnotationConfigApplicationContext;
import org.springframework.context.annotation.Bean;
import org.springframework.integration.dsl.IntegrationFlows;
import org.springframework.messaging.Message;
@@ -52,54 +52,75 @@ public class StreamEmitterBasicTests {
@Test
public void testFluxReturnAndOutputMethodLevel() throws Exception {
AnnotationConfigApplicationContext context = new AnnotationConfigApplicationContext();
context.register(TestFluxReturnAndOutputMethodLevel.class);
context.refresh();
ConfigurableApplicationContext context = SpringApplication.run(TestFluxReturnAndOutputMethodLevel.class,
"--server.port=0",
"--spring.jmx.enabled=false",
"--spring.cloud.stream.bindings.input.contentType=text/plain",
"--spring.cloud.stream.bindings.output.contentType=text/plain");
receiveAndValidate(context);
context.close();
}
@Test
public void testVoidReturnAndOutputMethodParameter() throws Exception {
AnnotationConfigApplicationContext context = new AnnotationConfigApplicationContext();
context.register(TestVoidReturnAndOutputMethodParameter.class);
context.refresh();
ConfigurableApplicationContext context = SpringApplication.run(TestVoidReturnAndOutputMethodParameter.class,
"--server.port=0",
"--spring.jmx.enabled=false",
"--spring.cloud.stream.bindings.input.contentType=text/plain",
"--spring.cloud.stream.bindings.output.contentType=text/plain");
receiveAndValidate(context);
context.close();
}
@Test
public void testVoidReturnAndOutputAtMethodLevel() throws Exception {
AnnotationConfigApplicationContext context = new AnnotationConfigApplicationContext();
context.register(TestVoidReturnAndOutputAtMethodLevel.class);
context.refresh();
ConfigurableApplicationContext context = SpringApplication.run(TestVoidReturnAndOutputAtMethodLevel.class,
"--server.port=0",
"--spring.jmx.enabled=false",
"--spring.cloud.stream.bindings.input.contentType=text/plain",
"--spring.cloud.stream.bindings.output.contentType=text/plain");
receiveAndValidate(context);
context.close();
}
@Test
public void testVoidReturnAndMultipleOutputMethodParameters() throws Exception {
AnnotationConfigApplicationContext context = new AnnotationConfigApplicationContext();
context.register(TestVoidReturnAndMultipleOutputMethodParameters.class);
context.refresh();
ConfigurableApplicationContext context = SpringApplication.run(TestVoidReturnAndMultipleOutputMethodParameters.class,
"--server.port=0",
"--spring.jmx.enabled=false",
"--spring.cloud.stream.bindings.input.contentType=text/plain",
"--spring.cloud.stream.bindings.output.contentType=text/plain",
"--spring.cloud.stream.bindings.output1.contentType=text/plain",
"--spring.cloud.stream.bindings.output2.contentType=text/plain",
"--spring.cloud.stream.bindings.output3.contentType=text/plain");
receiveAndValidateMultipleOutputs(context);
context.close();
}
@Test
public void testMultipleStreamEmitterMethods() throws Exception {
AnnotationConfigApplicationContext context = new AnnotationConfigApplicationContext();
context.register(TestMultipleStreamEmitterMethods.class);
context.refresh();
ConfigurableApplicationContext context = SpringApplication.run(TestMultipleStreamEmitterMethods.class,
"--server.port=0",
"--spring.jmx.enabled=false",
"--spring.cloud.stream.bindings.input.contentType=text/plain",
"--spring.cloud.stream.bindings.output.contentType=text/plain",
"--spring.cloud.stream.bindings.output1.contentType=text/plain",
"--spring.cloud.stream.bindings.output2.contentType=text/plain",
"--spring.cloud.stream.bindings.output3.contentType=text/plain");
receiveAndValidateMultipleOutputs(context);
context.close();
}
@Test
public void testSameAppContextWithMultipleStreamEmitters() throws Exception {
AnnotationConfigApplicationContext context = new AnnotationConfigApplicationContext();
context.register(TestSameAppContextWithMultipleStreamEmitters.class);
context.refresh();
ConfigurableApplicationContext context = SpringApplication.run(TestSameAppContextWithMultipleStreamEmitters.class,
"--server.port=0",
"--spring.jmx.enabled=false",
"--spring.cloud.stream.bindings.input.contentType=text/plain",
"--spring.cloud.stream.bindings.output.contentType=text/plain",
"--spring.cloud.stream.bindings.output1.contentType=text/plain",
"--spring.cloud.stream.bindings.output2.contentType=text/plain",
"--spring.cloud.stream.bindings.output3.contentType=text/plain");
receiveAndValidateMultiStreamEmittersInSameContext(context);
context.close();
}
@@ -108,12 +129,12 @@ public class StreamEmitterBasicTests {
private static void receiveAndValidate(ConfigurableApplicationContext context) throws InterruptedException {
Source source = context.getBean(Source.class);
MessageCollector messageCollector = context.getBean(MessageCollector.class);
List<String> messages = new ArrayList<>();
List<byte[]> messages = new ArrayList<>();
for (int i = 0; i < 1000; i++) {
messages.add((String) messageCollector.forChannel(source.output()).poll(5000, TimeUnit.MILLISECONDS).getPayload());
messages.add((byte[]) messageCollector.forChannel(source.output()).poll(5000, TimeUnit.MILLISECONDS).getPayload());
}
for (int i = 0; i < 1000; i++) {
assertThat(messages.get(i)).isEqualTo("HELLO WORLD!!" + i);
assertThat(new String(messages.get(i))).isEqualTo("HELLO WORLD!!" + i);
}
}
@@ -121,7 +142,7 @@ public class StreamEmitterBasicTests {
private static void receiveAndValidateMultipleOutputs(ConfigurableApplicationContext context) throws InterruptedException {
TestMultiOutboundChannels source = context.getBean(TestMultiOutboundChannels.class);
MessageCollector messageCollector = context.getBean(MessageCollector.class);
List<String> messages = new ArrayList<>();
List<byte[]> messages = new ArrayList<>();
assertMessages(source.output1(), messageCollector, messages);
messages.clear();
assertMessages(source.output2(), messageCollector, messages);
@@ -135,37 +156,37 @@ public class StreamEmitterBasicTests {
TestMultiOutboundChannels source1 = context1.getBean(TestMultiOutboundChannels.class);
MessageCollector messageCollector = context1.getBean(MessageCollector.class);
List<String> messages = new ArrayList<>();
List<byte[]> messages = new ArrayList<>();
assertMessagesX(source1.output1(), messageCollector, messages);
messages.clear();
assertMessagesY(source1.output2(), messageCollector, messages);
messages.clear();
}
private static void assertMessages(MessageChannel channel, MessageCollector messageCollector, List<String> messages) throws InterruptedException {
private static void assertMessages(MessageChannel channel, MessageCollector messageCollector, List<byte[]> messages) throws InterruptedException {
for (int i = 0; i < 1000; i++) {
messages.add((String) messageCollector.forChannel(channel).poll(5000, TimeUnit.MILLISECONDS).getPayload());
messages.add((byte[]) messageCollector.forChannel(channel).poll(5000, TimeUnit.MILLISECONDS).getPayload());
}
for (int i = 0; i < 1000; i++) {
assertThat(messages.get(i)).isEqualTo("Hello World!!" + i);
assertThat(new String(messages.get(i))).isEqualTo("Hello World!!" + i);
}
}
private static void assertMessagesX(MessageChannel channel, MessageCollector messageCollector, List<String> messages) throws InterruptedException {
private static void assertMessagesX(MessageChannel channel, MessageCollector messageCollector, List<byte[]> messages) throws InterruptedException {
for (int i = 0; i < 1000; i++) {
messages.add((String) messageCollector.forChannel(channel).poll(5000, TimeUnit.MILLISECONDS).getPayload());
messages.add((byte[]) messageCollector.forChannel(channel).poll(5000, TimeUnit.MILLISECONDS).getPayload());
}
for (int i = 0; i < 1000; i++) {
assertThat(messages.get(i)).isEqualTo("Hello World!!" + i);
assertThat(new String(messages.get(i))).isEqualTo("Hello World!!" + i);
}
}
private static void assertMessagesY(MessageChannel channel, MessageCollector messageCollector, List<String> messages) throws InterruptedException {
private static void assertMessagesY(MessageChannel channel, MessageCollector messageCollector, List<byte[]> messages) throws InterruptedException {
for (int i = 0; i < 1000; i++) {
messages.add((String) messageCollector.forChannel(channel).poll(5000, TimeUnit.MILLISECONDS).getPayload());
messages.add((byte[]) messageCollector.forChannel(channel).poll(5000, TimeUnit.MILLISECONDS).getPayload());
}
for (int i = 0; i < 1000; i++) {
assertThat(messages.get(i)).isEqualTo("Hello FooBar!!" + i);
assertThat(new String(messages.get(i))).isEqualTo("Hello FooBar!!" + i);
}
}

View File

@@ -40,25 +40,32 @@ import static org.springframework.cloud.stream.binding.StreamListenerErrorMessag
/**
* @author Ilayaperumal Gopinathan
* @author Vinicius Carvalho
*/
@SuppressWarnings("unchecked")
public class StreamListenerGenericFluxInputOutputArgsWithMessageTests {
@SuppressWarnings("unchecked")
private static void sendMessageAndValidate(ConfigurableApplicationContext context) throws InterruptedException {
private static void sendMessageAndValidate(ConfigurableApplicationContext context)
throws InterruptedException {
Processor processor = context.getBean(Processor.class);
String sentPayload = "hello " + UUID.randomUUID().toString();
processor.input().send(MessageBuilder.withPayload(sentPayload).setHeader("contentType", "text/plain").build());
processor.input().send(MessageBuilder.withPayload(sentPayload)
.setHeader("contentType", "text/plain").build());
MessageCollector messageCollector = context.getBean(MessageCollector.class);
Message<?> result = messageCollector.forChannel(processor.output()).poll(1000, TimeUnit.MILLISECONDS);
Message<byte[]> result = (Message<byte[]>) messageCollector.forChannel(processor.output()).poll(1000,
TimeUnit.MILLISECONDS);
assertThat(result).isNotNull();
assertThat(result.getPayload()).isEqualTo(sentPayload.toUpperCase());
assertThat(result.getPayload()).isEqualTo(sentPayload.toUpperCase().getBytes());
}
@Test
public void testGenericFluxInputOutputArgsWithMessage() throws Exception {
ConfigurableApplicationContext context = SpringApplication
.run(TestGenericStringFluxInputOutputArgsWithMessageImpl1.class, "--server.port=0");
ConfigurableApplicationContext context = SpringApplication.run(
TestGenericStringFluxInputOutputArgsWithMessageImpl1.class,
"--server.port=0", "--spring.jmx.enabled=false",
"--spring.cloud.stream.bindings.input.contentType=text/plain",
"--spring.cloud.stream.bindings.output.contentType=text/plain");
sendMessageAndValidate(context);
context.close();
}
@@ -66,11 +73,16 @@ public class StreamListenerGenericFluxInputOutputArgsWithMessageTests {
@Test
public void testInvalidInputValueWithOutputMethodParameters() {
try {
SpringApplication.run(TestGenericStringFluxInputOutputArgsWithMessageImpl2.class, "--server.port=0");
SpringApplication.run(
TestGenericStringFluxInputOutputArgsWithMessageImpl2.class,
"--server.port=0", "--spring.jmx.enabled=false",
"--spring.cloud.stream.bindings.input.contentType=text/plain",
"--spring.cloud.stream.bindings.output.contentType=text/plain");
fail("Expected exception: " + INVALID_INPUT_VALUE_WITH_OUTPUT_METHOD_PARAM);
}
catch (Exception e) {
assertThat(e.getMessage()).contains(INVALID_INPUT_VALUE_WITH_OUTPUT_METHOD_PARAM);
assertThat(e.getMessage())
.contains(INVALID_INPUT_VALUE_WITH_OUTPUT_METHOD_PARAM);
}
}
@@ -89,7 +101,8 @@ public class StreamListenerGenericFluxInputOutputArgsWithMessageTests {
@StreamListener
public void receive(@Input(Processor.INPUT) Flux<A> input,
@Output(Processor.OUTPUT) FluxSender output) {
output.send(input.map(m -> MessageBuilder.withPayload((A) m.toString().toUpperCase()).build()));
output.send(input.map(m -> MessageBuilder
.withPayload((A) m.toString().toUpperCase()).build()));
}
}
@@ -98,9 +111,9 @@ public class StreamListenerGenericFluxInputOutputArgsWithMessageTests {
public static class TestGenericFluxInputOutputArgsWithMessage2<A> {
@StreamListener(Processor.INPUT)
public void receive(Flux<A> input,
@Output(Processor.OUTPUT) FluxSender output) {
output.send(input.map(m -> MessageBuilder.withPayload((A) m.toString().toUpperCase()).build()));
public void receive(Flux<A> input, @Output(Processor.OUTPUT) FluxSender output) {
output.send(input.map(m -> MessageBuilder
.withPayload((A) m.toString().toUpperCase()).build()));
}
}
}

View File

@@ -65,14 +65,17 @@ public class StreamListenerReactiveInputOutputArgsTests {
String sentPayload = "hello " + UUID.randomUUID().toString();
processor.input().send(MessageBuilder.withPayload(sentPayload).setHeader("contentType", "text/plain").build());
MessageCollector messageCollector = context.getBean(MessageCollector.class);
Message<?> result = messageCollector.forChannel(processor.output()).poll(1000, TimeUnit.MILLISECONDS);
Message<byte[]> result = (Message<byte[]>) messageCollector.forChannel(processor.output()).poll(1000, TimeUnit.MILLISECONDS);
assertThat(result).isNotNull();
assertThat(result.getPayload()).isEqualTo(sentPayload.toUpperCase());
assertThat(result.getPayload()).isEqualTo(sentPayload.toUpperCase().getBytes());
}
@Test
public void testInputOutputArgs() throws Exception {
ConfigurableApplicationContext context = SpringApplication.run(this.configClass, "--server.port=0");
ConfigurableApplicationContext context = SpringApplication.run(this.configClass, "--server.port=0",
"--spring.jmx.enabled=false",
"--spring.cloud.stream.bindings.input.contentType=text/plain",
"--spring.cloud.stream.bindings.output.contentType=text/plain");
sendMessageAndValidate(context);
context.close();
}

View File

@@ -66,14 +66,17 @@ public class StreamListenerReactiveInputOutputArgsWithMessageTests {
String sentPayload = "hello " + UUID.randomUUID().toString();
processor.input().send(MessageBuilder.withPayload(sentPayload).setHeader("contentType", "text/plain").build());
MessageCollector messageCollector = context.getBean(MessageCollector.class);
Message<?> result = messageCollector.forChannel(processor.output()).poll(1000, TimeUnit.MILLISECONDS);
Message<byte[]> result = (Message<byte[]>) messageCollector.forChannel(processor.output()).poll(1000, TimeUnit.MILLISECONDS);
assertThat(result).isNotNull();
assertThat(result.getPayload()).isEqualTo(sentPayload.toUpperCase());
assertThat(result.getPayload()).isEqualTo(sentPayload.toUpperCase().getBytes());
}
@Test
public void testInputOutputArgs() throws Exception {
ConfigurableApplicationContext context = SpringApplication.run(this.configClass, "--server.port=0");
ConfigurableApplicationContext context = SpringApplication.run(this.configClass, "--server.port=0",
"--spring.jmx.enabled=false",
"--spring.cloud.stream.bindings.input.contentType=text/plain",
"--spring.cloud.stream.bindings.output.contentType=text/plain");
sendMessageAndValidate(context);
context.close();
}

View File

@@ -66,9 +66,9 @@ public class StreamListenerReactiveInputOutputArgsWithSenderAndFailureTests {
String sentPayload = "hello " + UUID.randomUUID().toString();
processor.input().send(MessageBuilder.withPayload(sentPayload).setHeader("contentType", "text/plain").build());
MessageCollector messageCollector = context.getBean(MessageCollector.class);
Message<?> result = messageCollector.forChannel(processor.output()).poll(1000, TimeUnit.MILLISECONDS);
Message<byte[]> result = (Message<byte[]>) messageCollector.forChannel(processor.output()).poll(1000, TimeUnit.MILLISECONDS);
assertThat(result).isNotNull();
assertThat(result.getPayload()).isEqualTo(sentPayload.toUpperCase());
assertThat(result.getPayload()).isEqualTo(sentPayload.toUpperCase().getBytes());
}
private static void sendFailingMessage(ConfigurableApplicationContext context) throws InterruptedException {
@@ -79,7 +79,10 @@ public class StreamListenerReactiveInputOutputArgsWithSenderAndFailureTests {
@Test
public void testInputOutputArgsWithFluxSenderAndFailure() throws Exception {
ConfigurableApplicationContext context = SpringApplication.run(this.configClass, "--server.port=0");
ConfigurableApplicationContext context = SpringApplication.run(this.configClass, "--server.port=0",
"--spring.jmx.enabled=false",
"--spring.cloud.stream.bindings.input.contentType=text/plain",
"--spring.cloud.stream.bindings.output.contentType=text/plain");
sendMessageAndValidate(context);
sendFailingMessage(context);
sendMessageAndValidate(context);

View File

@@ -66,15 +66,18 @@ public class StreamListenerReactiveInputOutputArgsWithSenderTests {
String sentPayload = "hello " + UUID.randomUUID().toString();
processor.input().send(MessageBuilder.withPayload(sentPayload).setHeader("contentType", "text/plain").build());
MessageCollector messageCollector = context.getBean(MessageCollector.class);
Message<?> result = messageCollector.forChannel(processor.output()).poll(1000, TimeUnit.MILLISECONDS);
Message<byte[]> result = (Message<byte[]>) messageCollector.forChannel(processor.output()).poll(1000, TimeUnit.MILLISECONDS);
assertThat(result).isNotNull();
assertThat(result.getPayload()).isEqualTo(sentPayload.toUpperCase());
assertThat(result.getPayload()).isEqualTo(sentPayload.toUpperCase().getBytes());
}
@Test
public void testInputOutputArgsWithFluxSender() throws Exception {
ConfigurableApplicationContext context = SpringApplication.run(this.configClass,
"--server.port=0");
"--server.port=0",
"--spring.jmx.enabled=false",
"--spring.cloud.stream.bindings.input.contentType=text/plain",
"--spring.cloud.stream.bindings.output.contentType=text/plain");
// send multiple message
sendMessageAndValidate(context);
sendMessageAndValidate(context);

View File

@@ -52,7 +52,10 @@ public class StreamListenerReactiveMethodTests {
@Test
public void testRxJava1InvalidInputValueWithOutputMethodParameters() {
try {
SpringApplication.run(RxJava1TestInputOutputArgs.class, "--server.port=0");
SpringApplication.run(RxJava1TestInputOutputArgs.class, "--server.port=0",
"--spring.jmx.enabled=false",
"--spring.cloud.stream.bindings.input.contentType=text/plain",
"--spring.cloud.stream.bindings.output.contentType=text/plain");
fail("IllegalArgumentException should have been thrown");
}
catch (Exception e) {
@@ -63,7 +66,10 @@ public class StreamListenerReactiveMethodTests {
@Test
public void testMethodReturnTypeWithNoOutboundSpecified() {
try {
SpringApplication.run(ReactorTestReturn5.class, "--server.port=0");
SpringApplication.run(ReactorTestReturn5.class, "--server.port=0",
"--spring.jmx.enabled=false",
"--spring.cloud.stream.bindings.input.contentType=text/plain",
"--spring.cloud.stream.bindings.output.contentType=text/plain");
fail("Exception expected: " + RETURN_TYPE_NO_OUTBOUND_SPECIFIED);
}
catch (Exception e) {

View File

@@ -69,14 +69,17 @@ public class StreamListenerReactiveMethodWithReturnTypeTests {
String sentPayload = "hello " + UUID.randomUUID().toString();
processor.input().send(MessageBuilder.withPayload(sentPayload).setHeader("contentType", "text/plain").build());
MessageCollector messageCollector = context.getBean(MessageCollector.class);
Message<?> result = messageCollector.forChannel(processor.output()).poll(1000, TimeUnit.MILLISECONDS);
Message<byte[]> result = (Message<byte[]>) messageCollector.forChannel(processor.output()).poll(1000, TimeUnit.MILLISECONDS);
assertThat(result).isNotNull();
assertThat(result.getPayload()).isEqualTo(sentPayload.toUpperCase());
assertThat(result.getPayload()).isEqualTo(sentPayload.toUpperCase().getBytes());
}
@Test
public void testReturn() throws Exception {
ConfigurableApplicationContext context = SpringApplication.run(this.configClass, "--server.port=0");
ConfigurableApplicationContext context = SpringApplication.run(this.configClass, "--server.port=0",
"--spring.jmx.enabled=false",
"--spring.cloud.stream.bindings.input.contentType=text/plain",
"--spring.cloud.stream.bindings.output.contentType=text/plain");
sendMessageAndValidate(context);
sendMessageAndValidate(context);
sendMessageAndValidate(context);

View File

@@ -70,9 +70,9 @@ public class StreamListenerReactiveReturnWithFailureTests {
String sentPayload = "hello " + UUID.randomUUID().toString();
processor.input().send(MessageBuilder.withPayload(sentPayload).setHeader("contentType", "text/plain").build());
MessageCollector messageCollector = context.getBean(MessageCollector.class);
Message<?> result = messageCollector.forChannel(processor.output()).poll(1000, TimeUnit.MILLISECONDS);
Message<byte[]> result = (Message<byte[]>) messageCollector.forChannel(processor.output()).poll(1000, TimeUnit.MILLISECONDS);
assertThat(result).isNotNull();
assertThat(result.getPayload()).isEqualTo(sentPayload.toUpperCase());
assertThat(result.getPayload()).isEqualTo(sentPayload.toUpperCase().getBytes());
}
private static void sendFailingMessage(ConfigurableApplicationContext context) throws InterruptedException {
@@ -83,7 +83,10 @@ public class StreamListenerReactiveReturnWithFailureTests {
@Test
public void testReturnWithFailure() throws Exception {
ConfigurableApplicationContext context = SpringApplication.run(this.configClass, "--server.port=0");
ConfigurableApplicationContext context = SpringApplication.run(this.configClass, "--server.port=0",
"--spring.jmx.enabled=false",
"--spring.cloud.stream.bindings.input.contentType=text/plain",
"--spring.cloud.stream.bindings.output.contentType=text/plain");
sendMessageAndValidate(context);
sendFailingMessage(context);
sendMessageAndValidate(context);

View File

@@ -70,14 +70,17 @@ public class StreamListenerReactiveReturnWithMessageTests {
String sentPayload = "hello " + UUID.randomUUID().toString();
processor.input().send(MessageBuilder.withPayload(sentPayload).setHeader("contentType", "text/plain").build());
MessageCollector messageCollector = context.getBean(MessageCollector.class);
Message<?> result = messageCollector.forChannel(processor.output()).poll(1000, TimeUnit.MILLISECONDS);
Message<byte[]> result = (Message<byte[]>) messageCollector.forChannel(processor.output()).poll(1000, TimeUnit.MILLISECONDS);
assertThat(result).isNotNull();
assertThat(result.getPayload()).isEqualTo(sentPayload.toUpperCase());
assertThat(result.getPayload()).isEqualTo(sentPayload.toUpperCase().getBytes());
}
@Test
public void testReturnWithMessage() throws Exception {
ConfigurableApplicationContext context = SpringApplication.run(this.configClass, "--server.port=0");
ConfigurableApplicationContext context = SpringApplication.run(this.configClass, "--server.port=0",
"--spring.jmx.enabled=false",
"--spring.cloud.stream.bindings.input.contentType=text/plain",
"--spring.cloud.stream.bindings.output.contentType=text/plain");
sendMessageAndValidate(context);
context.close();
}

View File

@@ -20,6 +20,9 @@ import java.util.Arrays;
import java.util.Collection;
import java.util.concurrent.TimeUnit;
import com.fasterxml.jackson.annotation.JsonCreator;
import com.fasterxml.jackson.annotation.JsonProperty;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.junit.runners.Parameterized;
@@ -50,6 +53,8 @@ public class StreamListenerReactiveReturnWithPojoTests {
private Class<?> configClass;
private ObjectMapper mapper = new ObjectMapper();
public StreamListenerReactiveReturnWithPojoTests(Class<?> configClass) {
this.configClass = configClass;
}
@@ -63,16 +68,18 @@ public class StreamListenerReactiveReturnWithPojoTests {
@Test
public void testReturnWithPojo() throws Exception {
ConfigurableApplicationContext context = SpringApplication.run(this.configClass, "--server.port=0");
ConfigurableApplicationContext context = SpringApplication.run(this.configClass, "--server.port=0",
"--spring.jmx.enabled=false");
@SuppressWarnings("unchecked")
Processor processor = context.getBean(Processor.class);
processor.input().send(MessageBuilder.withPayload("{\"message\":\"helloPojo\"}")
.setHeader("contentType", "application/json").build());
MessageCollector messageCollector = context.getBean(MessageCollector.class);
Message<?> result = messageCollector.forChannel(processor.output()).poll(1000, TimeUnit.MILLISECONDS);
Message<byte[]> result = (Message<byte[]>) messageCollector.forChannel(processor.output()).poll(1000, TimeUnit.MILLISECONDS);
assertThat(result).isNotNull();
assertThat(result.getPayload()).isInstanceOf(BarPojo.class);
assertThat(((BarPojo) result.getPayload()).getBarMessage()).isEqualTo("helloPojo");
assertThat(result.getPayload()).isInstanceOf(byte[].class);
BarPojo barPojo = mapper.readValue(result.getPayload(),BarPojo.class);
assertThat(barPojo.getBarMessage()).isEqualTo("helloPojo");
context.close();
}
@@ -175,7 +182,8 @@ public class StreamListenerReactiveReturnWithPojoTests {
private String barMessage;
public BarPojo(String barMessage) {
@JsonCreator
public BarPojo(@JsonProperty("barMessage") String barMessage) {
this.barMessage = barMessage;
}

View File

@@ -51,15 +51,15 @@ public class StreamListenerWildCardFluxInputOutputArgsWithMessageTests {
String sentPayload = "hello " + UUID.randomUUID().toString();
processor.input().send(MessageBuilder.withPayload(sentPayload).setHeader("contentType", "text/plain").build());
MessageCollector messageCollector = context.getBean(MessageCollector.class);
Message<?> result = messageCollector.forChannel(processor.output()).poll(1000, TimeUnit.MILLISECONDS);
Message<byte[]> result = (Message<byte[]>) messageCollector.forChannel(processor.output()).poll(1000, TimeUnit.MILLISECONDS);
assertThat(result).isNotNull();
assertThat(result.getPayload()).isEqualTo(sentPayload.toUpperCase());
assertThat(result.getPayload()).isEqualTo(sentPayload.toUpperCase().getBytes());
}
@Test
public void testWildCardFluxInputOutputArgsWithMessage() throws Exception {
ConfigurableApplicationContext context = SpringApplication
.run(TestWildCardFluxInputOutputArgsWithMessage1.class, "--server.port=0");
.run(TestWildCardFluxInputOutputArgsWithMessage1.class, "--server.port=0","--spring.cloud.stream.bindings.output.contentType=text/plain");
sendMessageAndValidate(context);
context.close();
}