From 0756b998bdb1233d4d39bc80c874924389dea5c4 Mon Sep 17 00:00:00 2001 From: Marius Bogoevici Date: Wed, 24 Aug 2016 14:50:13 -0400 Subject: [PATCH] Restore backwards compatibility with applications using input.contentType Fixes #633 Fixing checkstyle errors --- .../DeserializeJSONToJavaTypeTests.java | 93 +++++++++++++++++++ .../config/fooprocesor/foo-sink.properties | 2 + .../binding/MessageConverterConfigurer.java | 2 +- .../CompositeMessageConverterFactory.java | 1 + .../converter/JsonUnmarshallingConverter.java | 86 +++++++++++++++++ 5 files changed, 183 insertions(+), 1 deletion(-) create mode 100644 spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/DeserializeJSONToJavaTypeTests.java create mode 100644 spring-cloud-stream-integration-tests/src/test/resources/org/springframework/cloud/stream/config/fooprocesor/foo-sink.properties create mode 100644 spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/JsonUnmarshallingConverter.java diff --git a/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/DeserializeJSONToJavaTypeTests.java b/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/DeserializeJSONToJavaTypeTests.java new file mode 100644 index 000000000..f01e224a0 --- /dev/null +++ b/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/DeserializeJSONToJavaTypeTests.java @@ -0,0 +1,93 @@ +/* + * Copyright 2016 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 + * + * http://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.stream.config; + +import java.util.List; +import java.util.concurrent.TimeUnit; + +import org.junit.Test; +import org.junit.runner.RunWith; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.autoconfigure.EnableAutoConfiguration; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.cloud.stream.annotation.EnableBinding; +import org.springframework.cloud.stream.binder.BinderFactory; +import org.springframework.cloud.stream.messaging.Processor; +import org.springframework.cloud.stream.test.binder.TestSupportBinder; +import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.PropertySource; +import org.springframework.integration.annotation.ServiceActivator; +import org.springframework.integration.support.MessageBuilder; +import org.springframework.messaging.Message; +import org.springframework.messaging.converter.MessageConverter; +import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; + +import static org.assertj.core.api.Assertions.assertThat; + +/** + * @author Marius Bogoevici + */ +@RunWith(SpringJUnit4ClassRunner.class) +@SpringBootTest(classes = DeserializeJSONToJavaTypeTests.FooProcessor.class) +public class DeserializeJSONToJavaTypeTests { + + @Autowired + private Processor testProcessor; + + @Autowired + private BinderFactory binderFactory; + + @Autowired + private List customMessageConverters; + + @Test + public void testMessageDeserialized() throws Exception { + testProcessor.input().send(MessageBuilder.withPayload("{\"name\":\"Bar\"}").setHeader("contentType", "application/json").build()); + @SuppressWarnings("unchecked") + Message received = ((TestSupportBinder) binderFactory.getBinder(null)) + .messageCollector().forChannel(testProcessor.output()).poll(1, TimeUnit.SECONDS); + assertThat(received).isNotNull(); + assertThat(received.getPayload()).isInstanceOf(Foo.class); + assertThat((Foo) received.getPayload()).hasFieldOrPropertyWithValue("name", "Bar"); + } + + @EnableBinding(Processor.class) + @EnableAutoConfiguration + @PropertySource("classpath:/org/springframework/cloud/stream/config/fooprocesor/foo-sink.properties") + @Configuration + public static class FooProcessor { + + @ServiceActivator(inputChannel = "input", outputChannel = "output") + public Foo consume(Foo foo) { + return foo; + } + } + + public static class Foo { + + private String name; + + public String getName() { + return name; + } + + public void setName(String name) { + this.name = name; + } + } +} diff --git a/spring-cloud-stream-integration-tests/src/test/resources/org/springframework/cloud/stream/config/fooprocesor/foo-sink.properties b/spring-cloud-stream-integration-tests/src/test/resources/org/springframework/cloud/stream/config/fooprocesor/foo-sink.properties new file mode 100644 index 000000000..dea469faa --- /dev/null +++ b/spring-cloud-stream-integration-tests/src/test/resources/org/springframework/cloud/stream/config/fooprocesor/foo-sink.properties @@ -0,0 +1,2 @@ +spring.cloud.stream.bindings.input.destination=foo-input +spring.cloud.stream.bindings.input.content-type=application/x-java-object;type=org.springframework.cloud.stream.config.DeserializeJSONToJavaTypeTests.Foo diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/MessageConverterConfigurer.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/MessageConverterConfigurer.java index e49eb8614..1faeeae71 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/MessageConverterConfigurer.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/MessageConverterConfigurer.java @@ -130,7 +130,7 @@ public class MessageConverterConfigurer implements MessageChannelConfigurer, Bea this.contentType = contentType; this.mimeType = MessageConverterUtils.getMimeType(contentType); this.input = input; - if (MessageConverterUtils.X_JAVA_OBJECT.equals(this.mimeType)) { + if (MessageConverterUtils.X_JAVA_OBJECT.includes(this.mimeType)) { this.klazz = MessageConverterUtils .getJavaTypeForJavaObjectContentType(this.mimeType); diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/CompositeMessageConverterFactory.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/CompositeMessageConverterFactory.java index 4c12b5fca..d7055b1bb 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/CompositeMessageConverterFactory.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/CompositeMessageConverterFactory.java @@ -82,6 +82,7 @@ public class CompositeMessageConverterFactory { this.converters.add(new StringMessageConverter()); this.converters.add(new JavaSerializationMessageConverter()); + this.converters.add(new JsonUnmarshallingConverter(this.objectMapper)); } /** diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/JsonUnmarshallingConverter.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/JsonUnmarshallingConverter.java new file mode 100644 index 000000000..bfe5ae352 --- /dev/null +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/JsonUnmarshallingConverter.java @@ -0,0 +1,86 @@ +/* + * Copyright 2016 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 + * + * http://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.stream.converter; + +import java.io.IOException; + +import com.fasterxml.jackson.databind.ObjectMapper; + +import org.springframework.messaging.Message; +import org.springframework.messaging.MessageHeaders; +import org.springframework.messaging.converter.AbstractMessageConverter; +import org.springframework.messaging.converter.MessageConversionException; +import org.springframework.util.MimeType; +import org.springframework.util.MimeTypeUtils; + +/** + * Message converter providing backwards compatibility for applications using + * an Java type as input. + * + * @author Marius Bogoevici + */ +public class JsonUnmarshallingConverter extends AbstractMessageConverter { + + private final ObjectMapper objectMapper; + + protected JsonUnmarshallingConverter(ObjectMapper objectMapper) { + super(MessageConverterUtils.X_JAVA_OBJECT); + this.objectMapper = objectMapper != null ? objectMapper : new ObjectMapper(); + } + + @Override + protected boolean supports(Class aClass) { + return true; + } + + @Override + protected boolean canConvertFrom(Message message, Class targetClass) { + if ((message.getPayload() instanceof String) || (message.getPayload() instanceof byte[])) { + return true; + } + return canConvertFromBasedOnContentTypeHeader(message); + } + + private boolean canConvertFromBasedOnContentTypeHeader(Message message) { + Object contentTypeHeader = message.getHeaders().get(MessageHeaders.CONTENT_TYPE); + if (contentTypeHeader instanceof String) { + return MimeTypeUtils.APPLICATION_JSON.includes(MimeTypeUtils.parseMimeType((String) contentTypeHeader)); + } + else if (contentTypeHeader instanceof MimeType) { + return MimeTypeUtils.APPLICATION_JSON.includes((MimeType) contentTypeHeader); + } + else { + return contentTypeHeader == null; + } + } + + @Override + protected Object convertFromInternal(Message message, Class targetClass, Object conversionHint) { + Object payload = message.getPayload(); + try { + return payload instanceof byte[] ? objectMapper.readValue((byte[]) payload, targetClass) : objectMapper.readValue((String) payload, targetClass); + } + catch (IOException e) { + throw new MessageConversionException("Cannot parse payload ", e); + } + } + + @Override + protected Object convertToInternal(Object payload, MessageHeaders headers, Object conversionHint) { + return super.convertToInternal(payload, headers, conversionHint); + } +}