From dd5c8c2dec17ea3eb72ffb89390ee1f31caceffd Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Mon, 11 Jun 2018 14:59:00 -0400 Subject: [PATCH] GH-156: Fix outbound contentType Fixes https://github.com/spring-cloud/spring-cloud-stream-binder-rabbit/issues/156 The `SimpleMessageConverter` overwrote the `contentType` property with `application/octet-stream`. `contentType` set by the stream converter should win. - `setHeadersMappedLast(true)` - replace the converter with a simple `byte[]` pass-through converter on both sides. - if (for some reason) the stream converter produces something other than `byte[]` fall back to the `SimpleMessageConverter`. - remove the `afterPropertiesSet` since it's done by the abstract binder. **cherry-pick to 2.0.x** Resolves #157 --- .../rabbit/RabbitMessageChannelBinder.java | 36 ++++++++++++++++++- .../binder/rabbit/RabbitBinderTests.java | 6 ++++ 2 files changed, 41 insertions(+), 1 deletion(-) diff --git a/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java b/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java index 9dfb89541..2bceaa61d 100644 --- a/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java +++ b/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java @@ -40,6 +40,9 @@ import org.springframework.amqp.rabbit.retry.RejectAndDontRequeueRecoverer; import org.springframework.amqp.rabbit.retry.RepublishMessageRecoverer; import org.springframework.amqp.rabbit.support.DefaultMessagePropertiesConverter; import org.springframework.amqp.rabbit.support.MessagePropertiesConverter; +import org.springframework.amqp.support.converter.AbstractMessageConverter; +import org.springframework.amqp.support.converter.MessageConversionException; +import org.springframework.amqp.support.converter.SimpleMessageConverter; import org.springframework.amqp.support.postprocessor.DelegatingDecompressingPostProcessor; import org.springframework.amqp.support.postprocessor.GZipPostProcessor; import org.springframework.beans.factory.DisposableBean; @@ -107,6 +110,9 @@ public class RabbitMessageChannelBinder implements ExtendedPropertiesBinder, DisposableBean { + private static final SimplePassthroughMessageConverter passThoughConverter = + new SimplePassthroughMessageConverter(); + private static final AmqpMessageHeaderErrorMessageStrategy errorMessageStrategy = new AmqpMessageHeaderErrorMessageStrategy(); @@ -290,7 +296,7 @@ public class RabbitMessageChannelBinder endpoint.setConfirmCorrelationExpressionString("#root"); endpoint.setErrorMessageStrategy(new DefaultErrorMessageStrategy()); } - endpoint.afterPropertiesSet(); + endpoint.setHeadersMappedLast(true); return endpoint; } @@ -400,6 +406,7 @@ public class RabbitMessageChannelBinder adapter.setErrorMessageStrategy(errorMessageStrategy); adapter.setErrorChannel(errorInfrastructure.getErrorChannel()); } + adapter.setMessageConverter(passThoughConverter); return adapter; } @@ -590,6 +597,7 @@ public class RabbitMessageChannelBinder else { rabbitTemplate = new RabbitTemplate(); } + rabbitTemplate.setMessageConverter(passThoughConverter); rabbitTemplate.setChannelTransacted(properties.isTransacted()); rabbitTemplate.setConnectionFactory(this.connectionFactory); rabbitTemplate.setUsePublisherConnection(true); @@ -620,4 +628,30 @@ public class RabbitMessageChannelBinder return stringWriter.getBuffer().toString(); } + private static final class SimplePassthroughMessageConverter extends AbstractMessageConverter { + + private static final SimpleMessageConverter converter = new SimpleMessageConverter(); + + SimplePassthroughMessageConverter() { + super(); + } + + @Override + protected Message createMessage(Object object, MessageProperties messageProperties) { + if (object instanceof byte[]) { + return new Message((byte[]) object, messageProperties); + } + else { + // just for safety (backwards compatibility) + return converter.toMessage(object, messageProperties); + } + } + + @Override + public Object fromMessage(Message message) throws MessageConversionException { + return message.getBody(); + } + + } + } diff --git a/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/RabbitBinderTests.java b/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/RabbitBinderTests.java index 4e6293e90..0dbc62fb1 100644 --- a/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/RabbitBinderTests.java +++ b/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/RabbitBinderTests.java @@ -171,8 +171,14 @@ public class RabbitBinderTests extends DirectChannel moduleInputChannel = createBindableChannel("input", new BindingProperties()); Binding producerBinding = binder.bindProducer("bad.0", moduleOutputChannel, createProducerProperties()); + assertThat(TestUtils.getPropertyValue(producerBinding, "lifecycle.headersMappedLast", Boolean.class)) + .isTrue(); + assertThat(TestUtils.getPropertyValue(producerBinding, "lifecycle.amqpTemplate.messageConverter") + .getClass().getName()).contains("Passthrough"); Binding consumerBinding = binder.bindConsumer("bad.0", "test", moduleInputChannel, createConsumerProperties()); + assertThat(TestUtils.getPropertyValue(consumerBinding, "lifecycle.messageConverter") + .getClass().getName()).contains("Passthrough"); Message message = MessageBuilder.withPayload("bad".getBytes()).setHeader(MessageHeaders.CONTENT_TYPE, "foo/bar").build(); final CountDownLatch latch = new CountDownLatch(3); moduleInputChannel.subscribe(new MessageHandler() {