GH-1564 Delegated type-conversion to MessageConverters

The following is the summary of changes which essentially delegate all type conversion back to MessageConverter.
The only thing remains is the `BINDER_ORIGINAL_CONTENT_TYPE` logic to ensure backward compatibility

* `MessageConverterConfigurer` was brought pretty much back to the state it was before all those questionable type conversion changes
* `BinderFactoryConfiguration` configures  custom argument resolvers which will be removed as soon as https://jira.spring.io/browse/SPR-17503 is addressed.
* The two new argument resolvers, defer from their original counterparts in that they change the order of type assertion ensuring that, for example, byte[] does not match Object and would have to be sent to MessageConverter for possible conversion. These two resolvers will be removed once https://jira.spring.io/browse/SPR-17503 is addressed.

Resolves #1564
Resolves #1565
This commit is contained in:
Oleg Zhurakousky
2018-12-17 21:50:56 +01:00
parent d5aa0906f4
commit a731021b43
8 changed files with 311 additions and 61 deletions

View File

@@ -8,7 +8,7 @@
<parent>
<groupId>org.springframework.cloud</groupId>
<artifactId>spring-cloud-build</artifactId>
<version>2.1.0.RC3</version>
<version>2.1.0.BUILD-SNAPSHOT</version>
<relativePath/>
</parent>
<scm>
@@ -24,7 +24,7 @@
<reactor.version>Californium-RELEASE</reactor.version>
<kryo-shaded.version>3.0.3</kryo-shaded.version>
<objenesis.version>2.1</objenesis.version>
<spring-cloud-function.version>2.0.0.RC2</spring-cloud-function.version>
<spring-cloud-function.version>2.0.0.BUILD-SNAPSHOT</spring-cloud-function.version>
</properties>
<dependencyManagement>
<dependencies>

View File

@@ -181,10 +181,10 @@ public abstract class AbstractBinderTests<B extends AbstractTestBinder<? extends
binderBindUnbindLatency();
CountDownLatch latch = new CountDownLatch(1);
AtomicReference<Message<String>> inboundMessageRef = new AtomicReference<Message<String>>();
AtomicReference<Message<byte[]>> inboundMessageRef = new AtomicReference<Message<byte[]>>();
moduleInputChannel.subscribe(message1 -> {
try {
inboundMessageRef.set((Message<String>) message1);
inboundMessageRef.set((Message<byte[]>) message1);
}
finally {
latch.countDown();
@@ -194,7 +194,7 @@ public abstract class AbstractBinderTests<B extends AbstractTestBinder<? extends
moduleOutputChannel.send(message);
Assert.isTrue(latch.await(5, TimeUnit.SECONDS), "Failed to receive message");
assertThat(inboundMessageRef.get().getPayload()).isEqualTo("foo");
assertThat(inboundMessageRef.get().getPayload()).isEqualTo("foo".getBytes());
assertThat(inboundMessageRef.get().getHeaders().get(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE)).isNull();
assertThat(inboundMessageRef.get().getHeaders().get(MessageHeaders.CONTENT_TYPE).toString()).isEqualTo("text/plain");
producerBinding.unbind();
@@ -386,10 +386,10 @@ public abstract class AbstractBinderTests<B extends AbstractTestBinder<? extends
.setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.TEXT_PLAIN).build();
moduleOutputChannel.send(message);
CountDownLatch latch = new CountDownLatch(1);
AtomicReference<Message<String>> inboundMessageRef = new AtomicReference<Message<String>>();
AtomicReference<Message<byte[]>> inboundMessageRef = new AtomicReference<Message<byte[]>>();
moduleInputChannel.subscribe(message1 -> {
try {
inboundMessageRef.set((Message<String>) message1);
inboundMessageRef.set((Message<byte[]>) message1);
}
finally {
latch.countDown();
@@ -399,7 +399,7 @@ public abstract class AbstractBinderTests<B extends AbstractTestBinder<? extends
moduleOutputChannel.send(message);
Assert.isTrue(latch.await(5, TimeUnit.SECONDS), "Failed to receive message");
assertThat(inboundMessageRef.get()).isNotNull();
assertThat(inboundMessageRef.get().getPayload()).isEqualTo("foo");
assertThat(inboundMessageRef.get().getPayload()).isEqualTo("foo".getBytes());
assertThat(inboundMessageRef.get().getHeaders().get(MessageHeaders.CONTENT_TYPE).toString())
.isEqualTo(MimeTypeUtils.TEXT_PLAIN_VALUE);
producerBinding.unbind();

View File

@@ -17,7 +17,6 @@
package org.springframework.cloud.stream.binding;
import java.lang.reflect.Field;
import java.nio.charset.StandardCharsets;
import java.util.Map;
import org.apache.commons.logging.Log;
@@ -53,7 +52,6 @@ import org.springframework.messaging.converter.MessageConverter;
import org.springframework.messaging.handler.invocation.InvocableHandlerMethod;
import org.springframework.messaging.support.ChannelInterceptor;
import org.springframework.messaging.support.ErrorMessage;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.util.Assert;
import org.springframework.util.CollectionUtils;
import org.springframework.util.MimeType;
@@ -284,10 +282,6 @@ public class MessageConverterConfigurer implements MessageChannelAndSourceConfig
headersMap.put(MessageHeaders.CONTENT_TYPE, MimeType.valueOf((String)message.getHeaders().get(MessageHeaders.CONTENT_TYPE)));
}
if (message.getHeaders().get(MessageHeaders.CONTENT_TYPE).toString().startsWith("text") && message.getPayload() instanceof byte[]) {
message = MessageBuilder.withPayload(new String((byte[])message.getPayload(), StandardCharsets.UTF_8)).copyHeaders(message.getHeaders()).build();
}
return message;
}
}
@@ -310,19 +304,11 @@ public class MessageConverterConfigurer implements MessageChannelAndSourceConfig
@Override
public Message<?> 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<String, Object> headersMap = (Map<String, Object>) 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<String, Object> headersMap = (Map<String, Object>) 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;

View File

@@ -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<HandlerMethodArgumentResolver> 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;
}
}

View File

@@ -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;
}
}

View File

@@ -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;
}
}

View File

@@ -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());
}
}

View File

@@ -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<byte[]> 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)