diff --git a/pom.xml b/pom.xml
index 008dc5c59..98d15a885 100644
--- a/pom.xml
+++ b/pom.xml
@@ -8,7 +8,7 @@
org.springframework.cloud
spring-cloud-build
- 2.1.0.RC3
+ 2.1.0.BUILD-SNAPSHOT
@@ -24,7 +24,7 @@
Californium-RELEASE
3.0.3
2.1
- 2.0.0.RC2
+ 2.0.0.BUILD-SNAPSHOT
diff --git a/spring-cloud-stream-binder-test/src/main/java/org/springframework/cloud/stream/binder/AbstractBinderTests.java b/spring-cloud-stream-binder-test/src/main/java/org/springframework/cloud/stream/binder/AbstractBinderTests.java
index d9132d218..2d02e4650 100644
--- a/spring-cloud-stream-binder-test/src/main/java/org/springframework/cloud/stream/binder/AbstractBinderTests.java
+++ b/spring-cloud-stream-binder-test/src/main/java/org/springframework/cloud/stream/binder/AbstractBinderTests.java
@@ -181,10 +181,10 @@ public abstract class AbstractBinderTests> inboundMessageRef = new AtomicReference>();
+ AtomicReference> inboundMessageRef = new AtomicReference>();
moduleInputChannel.subscribe(message1 -> {
try {
- inboundMessageRef.set((Message) message1);
+ inboundMessageRef.set((Message) message1);
}
finally {
latch.countDown();
@@ -194,7 +194,7 @@ public abstract class AbstractBinderTests> inboundMessageRef = new AtomicReference>();
+ AtomicReference> inboundMessageRef = new AtomicReference>();
moduleInputChannel.subscribe(message1 -> {
try {
- inboundMessageRef.set((Message) message1);
+ inboundMessageRef.set((Message) message1);
}
finally {
latch.countDown();
@@ -399,7 +399,7 @@ public abstract class AbstractBinderTests doPreSend(Message> message, MessageChannel channel) {
- @SuppressWarnings("deprecation")
- boolean propagateOriginalContentType =
- MessageConverterConfigurer.this.bindingServiceProperties.isPropagateOriginalContentType();
-
// ===== 1.3 backward compatibility code part-1 ===
- String ct = null;
- String oct = null;
- if (propagateOriginalContentType) {
- oct = message.getHeaders().containsKey(MessageHeaders.CONTENT_TYPE) ? message.getHeaders().get(MessageHeaders.CONTENT_TYPE).toString() : null;
- ct = message.getPayload() instanceof String
+ String oct = message.getHeaders().containsKey(MessageHeaders.CONTENT_TYPE) ? message.getHeaders().get(MessageHeaders.CONTENT_TYPE).toString() : null;
+ String ct = message.getPayload() instanceof String
? ct = JavaClassMimeTypeUtils.mimeTypeFromObject(message.getPayload(), ObjectUtils.nullSafeToString(oct)).toString()
: oct;
- }
// ===== END 1.3 backward compatibility code part-1 ===
@@ -340,20 +326,12 @@ public class MessageConverterConfigurer implements MessageChannelAndSourceConfig
}
/// ===== 1.3 backward compatibility code part-2 ===
- if (propagateOriginalContentType) {
- if (ct != null && !ct.equals(oct) && oct != null) {
- @SuppressWarnings("unchecked")
- Map headersMap = (Map) ReflectionUtils.getField(MessageConverterConfigurer.this.headersField,
- outboundMessage.getHeaders());
- headersMap.put(MessageHeaders.CONTENT_TYPE, MimeType.valueOf(ct));
- headersMap.put(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE, MimeType.valueOf(oct));
- }
- }
- else {
+ if (ct != null && !ct.equals(oct) && oct != null) {
@SuppressWarnings("unchecked")
Map headersMap = (Map) ReflectionUtils.getField(MessageConverterConfigurer.this.headersField,
outboundMessage.getHeaders());
- headersMap.remove(MessageHeaders.CONTENT_TYPE);
+ headersMap.put(MessageHeaders.CONTENT_TYPE, MimeType.valueOf(ct));
+ headersMap.put(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE, MimeType.valueOf(oct));
}
// ===== END 1.3 backward compatibility code part-2 ===
return outboundMessage;
diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BinderFactoryConfiguration.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BinderFactoryConfiguration.java
index 28137d897..080f20daf 100644
--- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BinderFactoryConfiguration.java
+++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BinderFactoryConfiguration.java
@@ -22,6 +22,7 @@ import java.util.ArrayList;
import java.util.Collection;
import java.util.Enumeration;
import java.util.HashMap;
+import java.util.LinkedList;
import java.util.List;
import java.util.Map;
import java.util.Properties;
@@ -33,6 +34,7 @@ import org.springframework.beans.factory.BeanCreationException;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.beans.factory.config.BeanDefinition;
+import org.springframework.beans.factory.config.ConfigurableListableBeanFactory;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
import org.springframework.cloud.stream.binder.BinderType;
import org.springframework.cloud.stream.binder.BinderTypeRegistry;
@@ -55,7 +57,10 @@ import org.springframework.integration.context.IntegrationContextUtils;
import org.springframework.integration.handler.support.HandlerMethodArgumentResolversHolder;
import org.springframework.lang.Nullable;
import org.springframework.messaging.handler.annotation.support.DefaultMessageHandlerMethodFactory;
+import org.springframework.messaging.handler.annotation.support.HeaderMethodArgumentResolver;
+import org.springframework.messaging.handler.annotation.support.HeadersMethodArgumentResolver;
import org.springframework.messaging.handler.annotation.support.MessageHandlerMethodFactory;
+import org.springframework.messaging.handler.invocation.HandlerMethodArgumentResolver;
import org.springframework.util.ClassUtils;
import org.springframework.util.StringUtils;
import org.springframework.validation.Validator;
@@ -161,12 +166,33 @@ public class BinderFactoryConfiguration {
@Bean
public static MessageHandlerMethodFactory messageHandlerMethodFactory(CompositeMessageConverterFactory compositeMessageConverterFactory,
@Qualifier(IntegrationContextUtils.ARGUMENT_RESOLVERS_BEAN_NAME) HandlerMethodArgumentResolversHolder ahmar,
- @Nullable Validator validator) {
+ @Nullable Validator validator, ConfigurableListableBeanFactory clbf) {
+
DefaultMessageHandlerMethodFactory messageHandlerMethodFactory = new DefaultMessageHandlerMethodFactory();
messageHandlerMethodFactory.setMessageConverter(compositeMessageConverterFactory.getMessageConverterForAllRegistered());
- messageHandlerMethodFactory.setCustomArgumentResolvers(ahmar.getResolvers());
+
+ /*
+ * We essentially do the same thing as the DefaultMessageHandlerMethodFactory.initArgumentResolvers(..).
+ * We can't do it as custom resolvers for two reasons.
+ * 1. We would have two duplicate (compatible) resolvers, so they would need to be ordered properly
+ * to ensure these new resolvers take precedence.
+ * 2. DefaultMessageHandlerMethodFactory.initArgumentResolvers(..) puts MessageMethodArgumentResolver
+ * before custom converters thus not allowing an override which kind of proves #1.
+ *
+ * In all, all this will be obsolete once https://jira.spring.io/browse/SPR-17503 is addressed and we can fall back on
+ * core resolvers
+ */
+ List resolvers = new LinkedList<>();
+ resolvers.add(new SmartPayloadArgumentResolver(compositeMessageConverterFactory.getMessageConverterForAllRegistered(), validator));
+ resolvers.add(new SmartMessageMethodArgumentResolver(compositeMessageConverterFactory.getMessageConverterForAllRegistered()));
+ resolvers.add(new HeaderMethodArgumentResolver(null, clbf));
+ resolvers.add(new HeadersMethodArgumentResolver());
+ resolvers.addAll(ahmar.getResolvers());
+
+ messageHandlerMethodFactory.setArgumentResolvers(resolvers);
messageHandlerMethodFactory.setValidator(validator);
return messageHandlerMethodFactory;
}
+
}
diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BindingServiceProperties.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BindingServiceProperties.java
index 5a72880df..4294cffd1 100644
--- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BindingServiceProperties.java
+++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BindingServiceProperties.java
@@ -56,18 +56,6 @@ public class BindingServiceProperties implements ApplicationContextAware, Initia
private static final int DEFAULT_BINDING_RETRY_INTERVAL = 30;
- /**
- * Setting it to true ensures that the original content-type of the message is propagated
- * to the outgoing message as `originalContentType` header.
- *
- * This deprecated feature primarily exists for backward compatibility
- * and will not be supported in future versions.
- *
- * Default: true
- */
- @Deprecated
- private boolean propagateOriginalContentType = true;
-
/**
* The instance id of the application: a number from 0 to instanceCount-1.
* Used for partitioning and with Kafka.
@@ -295,15 +283,4 @@ public class BindingServiceProperties implements ApplicationContextAware, Initia
binder.bind("spring.cloud.stream.default", Bindable.ofInstance(bindingPropertiesTarget));
this.bindings.put(binding, bindingPropertiesTarget);
}
-
- @Deprecated
- public boolean isPropagateOriginalContentType() {
- return propagateOriginalContentType;
- }
-
- @Deprecated
- public void setPropagateOriginalContentType(boolean propagateOriginalContentType) {
- this.propagateOriginalContentType = propagateOriginalContentType;
- }
-
}
diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/SmartMessageMethodArgumentResolver.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/SmartMessageMethodArgumentResolver.java
new file mode 100644
index 000000000..be51a4b02
--- /dev/null
+++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/SmartMessageMethodArgumentResolver.java
@@ -0,0 +1,123 @@
+/*
+ * Copyright 2016-2018 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.lang.reflect.Type;
+
+import org.springframework.core.MethodParameter;
+import org.springframework.core.ResolvableType;
+import org.springframework.lang.Nullable;
+import org.springframework.messaging.Message;
+import org.springframework.messaging.converter.MessageConversionException;
+import org.springframework.messaging.converter.MessageConverter;
+import org.springframework.messaging.converter.SmartMessageConverter;
+import org.springframework.messaging.handler.annotation.support.MessageMethodArgumentResolver;
+import org.springframework.messaging.handler.annotation.support.MethodArgumentTypeMismatchException;
+import org.springframework.messaging.support.MessageBuilder;
+import org.springframework.util.ClassUtils;
+import org.springframework.util.StringUtils;
+
+/**
+ *
+ * @author Oleg Zhurakousky
+ *
+ * @deprecated will be removed once https://jira.spring.io/browse/SPR-17503 is addressed
+ */
+@Deprecated
+class SmartMessageMethodArgumentResolver extends MessageMethodArgumentResolver {
+
+ private final MessageConverter messageConverter;
+
+ SmartMessageMethodArgumentResolver() {
+ this(null);
+ }
+
+ /**
+ * Create a resolver instance with the given {@link MessageConverter}.
+ * @param converter the MessageConverter to use (may be {@code null})
+ * @since 4.3
+ */
+ SmartMessageMethodArgumentResolver(@Nullable MessageConverter converter) {
+ this.messageConverter = converter;
+ }
+
+ @Override
+ public Object resolveArgument(MethodParameter parameter, Message> message) throws Exception {
+ Class> targetMessageType = parameter.getParameterType();
+ Class> targetPayloadType = getPayloadType(parameter);
+
+ if (!targetMessageType.isAssignableFrom(message.getClass())) {
+ throw new MethodArgumentTypeMismatchException(message, parameter, "Actual message type '" +
+ ClassUtils.getDescriptiveType(message) + "' does not match expected type '" +
+ ClassUtils.getQualifiedName(targetMessageType) + "'");
+ }
+
+ Class> payloadClass = message.getPayload().getClass();
+
+ if (ClassUtils.isAssignable(payloadClass, targetPayloadType)) {
+ return message;
+ }
+ Object payload = message.getPayload();
+ if (isEmptyPayload(payload)) {
+ throw new MessageConversionException(message, "Cannot convert from actual payload type '" +
+ ClassUtils.getDescriptiveType(payload) + "' to expected payload type '" +
+ ClassUtils.getQualifiedName(targetPayloadType) + "' when payload is empty");
+ }
+
+ payload = convertPayload(message, parameter, targetPayloadType);
+ return MessageBuilder.createMessage(payload, message.getHeaders());
+ }
+
+ private Class> getPayloadType(MethodParameter parameter) {
+ Type genericParamType = parameter.getGenericParameterType();
+ ResolvableType resolvableType = ResolvableType.forType(genericParamType).as(Message.class);
+ return resolvableType.getGeneric().toClass();
+ }
+
+ protected boolean isEmptyPayload(@Nullable Object payload) {
+ if (payload == null) {
+ return true;
+ }
+ else if (payload instanceof byte[]) {
+ return ((byte[]) payload).length == 0;
+ }
+ else if (payload instanceof String) {
+ return !StringUtils.hasText((String) payload);
+ }
+ else {
+ return false;
+ }
+ }
+
+ private Object convertPayload(Message> message, MethodParameter parameter, Class> targetPayloadType) {
+ Object result = null;
+ if (this.messageConverter instanceof SmartMessageConverter) {
+ SmartMessageConverter smartConverter = (SmartMessageConverter) this.messageConverter;
+ result = smartConverter.fromMessage(message, targetPayloadType, parameter);
+ }
+ else if (this.messageConverter != null) {
+ result = this.messageConverter.fromMessage(message, targetPayloadType);
+ }
+
+ if (result == null) {
+ throw new MessageConversionException(message, "No converter found from actual payload type '" +
+ ClassUtils.getDescriptiveType(message.getPayload()) + "' to expected payload type '" +
+ ClassUtils.getQualifiedName(targetPayloadType) + "'");
+ }
+ return result;
+ }
+}
diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/SmartPayloadArgumentResolver.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/SmartPayloadArgumentResolver.java
new file mode 100644
index 000000000..35a56f787
--- /dev/null
+++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/SmartPayloadArgumentResolver.java
@@ -0,0 +1,119 @@
+/*
+ * Copyright 2016-2018 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 org.springframework.core.MethodParameter;
+import org.springframework.lang.Nullable;
+import org.springframework.messaging.Message;
+import org.springframework.messaging.converter.MessageConversionException;
+import org.springframework.messaging.converter.MessageConverter;
+import org.springframework.messaging.converter.SmartMessageConverter;
+import org.springframework.messaging.handler.annotation.Header;
+import org.springframework.messaging.handler.annotation.Headers;
+import org.springframework.messaging.handler.annotation.Payload;
+import org.springframework.messaging.handler.annotation.support.MethodArgumentNotValidException;
+import org.springframework.messaging.handler.annotation.support.PayloadArgumentResolver;
+import org.springframework.util.ClassUtils;
+import org.springframework.util.StringUtils;
+import org.springframework.validation.BeanPropertyBindingResult;
+import org.springframework.validation.BindingResult;
+import org.springframework.validation.ObjectError;
+import org.springframework.validation.Validator;
+
+/**
+ *
+ * @author Oleg Zhurakousky
+ *
+ * @deprecated will be removed once https://jira.spring.io/browse/SPR-17503 is addressed
+ */
+@Deprecated
+class SmartPayloadArgumentResolver extends PayloadArgumentResolver {
+
+ private final MessageConverter messageConverter;
+
+ SmartPayloadArgumentResolver(MessageConverter messageConverter) {
+ super(messageConverter);
+ this.messageConverter = messageConverter;
+ }
+
+ SmartPayloadArgumentResolver(MessageConverter messageConverter, Validator validator) {
+ super(messageConverter, validator, true);
+ this.messageConverter = messageConverter;
+ }
+
+ SmartPayloadArgumentResolver(MessageConverter messageConverter, Validator validator,
+ boolean useDefaultResolution) {
+ super(messageConverter, validator, useDefaultResolution);
+ this.messageConverter = messageConverter;
+ }
+
+ @Override
+ public boolean supportsParameter(MethodParameter parameter) {
+ return (!Message.class.isAssignableFrom(parameter.getParameterType())
+ && !parameter.hasParameterAnnotation(Header.class)
+ && !parameter.hasParameterAnnotation(Headers.class));
+ }
+
+ @Override
+ @Nullable
+ public Object resolveArgument(MethodParameter parameter, Message> message) throws Exception {
+ Payload ann = parameter.getParameterAnnotation(Payload.class);
+ if (ann != null && StringUtils.hasText(ann.expression())) {
+ throw new IllegalStateException("@Payload SpEL expressions not supported by this resolver");
+ }
+
+ Object payload = message.getPayload();
+ if (isEmptyPayload(payload)) {
+ if (ann == null || ann.required()) {
+ String paramName = getParameterName(parameter);
+ BindingResult bindingResult = new BeanPropertyBindingResult(payload, paramName);
+ bindingResult.addError(new ObjectError(paramName, "Payload value must not be empty"));
+ throw new MethodArgumentNotValidException(message, parameter, bindingResult);
+ }
+ else {
+ return null;
+ }
+ }
+
+ Class> targetClass = parameter.getParameterType();
+ Class> payloadClass = payload.getClass();
+ if (ClassUtils.isAssignable(payloadClass, targetClass)) {
+ validate(message, parameter, payload);
+ return payload;
+ }
+ else {
+ if (this.messageConverter instanceof SmartMessageConverter) {
+ SmartMessageConverter smartConverter = (SmartMessageConverter) this.messageConverter;
+ payload = smartConverter.fromMessage(message, targetClass, parameter);
+ }
+ else {
+ payload = this.messageConverter.fromMessage(message, targetClass);
+ }
+ if (payload == null) {
+ throw new MessageConversionException(message, "Cannot convert from [" +
+ payloadClass.getName() + "] to [" + targetClass.getName() + "] for " + message);
+ }
+ validate(message, parameter, payload);
+ return payload;
+ }
+ }
+
+ private String getParameterName(MethodParameter param) {
+ String paramName = param.getParameterName();
+ return (paramName != null ? paramName : "Arg " + param.getParameterIndex());
+ }
+}
diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/tck/ContentTypeTckTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/tck/ContentTypeTckTests.java
index dbea12ed0..298919574 100644
--- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/tck/ContentTypeTckTests.java
+++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/tck/ContentTypeTckTests.java
@@ -231,7 +231,21 @@ public class ContentTypeTckTests {
assertEquals(MimeTypeUtils.APPLICATION_JSON, outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE));
assertEquals(jsonPayload, new String(outputMessage.getPayload(), StandardCharsets.UTF_8));
}
-
+
+ @Test
+ public void typelessMessageToPojoInboundContentTypeBinding() {
+ ApplicationContext context = new SpringApplicationBuilder(TypelessMessageToPojoStreamListener.class)
+ .web(WebApplicationType.NONE)
+ .run("--spring.cloud.stream.bindings.input.contentType=text/plain", "--spring.jmx.enabled=false");
+ InputDestination source = context.getBean(InputDestination.class);
+ OutputDestination target = context.getBean(OutputDestination.class);
+ String jsonPayload = "{\"name\":\"oleg\"}";
+ source.send(new GenericMessage<>(jsonPayload.getBytes()));
+ Message outputMessage = target.receive();
+ assertEquals(MimeTypeUtils.APPLICATION_JSON, outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE));
+ assertEquals(jsonPayload, new String(outputMessage.getPayload(), StandardCharsets.UTF_8));
+ }
+
@Test
public void typelessToPojoWithTextHeaderContentTypeBinding() {
ApplicationContext context = new SpringApplicationBuilder(TypelessToPojoStreamListener.class)
@@ -614,6 +628,19 @@ public class ContentTypeTckTests {
return mapper.readValue((String)value, Person.class);
}
}
+
+ @EnableBinding(Processor.class)
+ @Import(TestChannelBinderConfiguration.class)
+ @EnableAutoConfiguration
+ public static class TypelessMessageToPojoStreamListener {
+ @StreamListener(Processor.INPUT)
+ @SendTo(Processor.OUTPUT)
+ public Person echo(Message> message) throws Exception {
+ ObjectMapper mapper = new ObjectMapper();
+ //assume it is string because CT is text/plain
+ return mapper.readValue((String)message.getPayload(), Person.class);
+ }
+ }
@EnableBinding(Processor.class)
@Import(TestChannelBinderConfiguration.class)