From f4f1785e208a461b7950776c31366a9e4898ac34 Mon Sep 17 00:00:00 2001 From: Vinicius Carvalho Date: Thu, 14 Sep 2017 15:19:05 -0400 Subject: [PATCH] ContentType revamp - Making tests compatible with new contentType behavior --- pom.xml | 5 --- spring-cloud-stream-binder-rabbit/pom.xml | 4 -- ...bbitMessageChannelBinderConfiguration.java | 8 +--- .../binder/rabbit/RabbitBinderTests.java | 42 +++++++------------ .../binder/rabbit/RabbitTestBinder.java | 2 - 5 files changed, 16 insertions(+), 45 deletions(-) diff --git a/pom.xml b/pom.xml index 98f074fb5..c6d6174e0 100644 --- a/pom.xml +++ b/pom.xml @@ -23,11 +23,6 @@ spring-cloud-stream ${spring-cloud-stream.version} - - org.springframework.cloud - spring-cloud-stream-codec - ${spring-cloud-stream.version} - org.springframework.cloud spring-cloud-stream-binder-rabbit diff --git a/spring-cloud-stream-binder-rabbit/pom.xml b/spring-cloud-stream-binder-rabbit/pom.xml index 382ef5de5..37e835194 100644 --- a/spring-cloud-stream-binder-rabbit/pom.xml +++ b/spring-cloud-stream-binder-rabbit/pom.xml @@ -28,10 +28,6 @@ org.springframework.cloud spring-cloud-stream - - org.springframework.cloud - spring-cloud-stream-codec - org.springframework.boot spring-boot-autoconfigure diff --git a/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/config/RabbitMessageChannelBinderConfiguration.java b/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/config/RabbitMessageChannelBinderConfiguration.java index b61f9b49a..6077fbb23 100644 --- a/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/config/RabbitMessageChannelBinderConfiguration.java +++ b/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/config/RabbitMessageChannelBinderConfiguration.java @@ -28,26 +28,23 @@ import org.springframework.cloud.stream.binder.rabbit.RabbitMessageChannelBinder import org.springframework.cloud.stream.binder.rabbit.properties.RabbitBinderConfigurationProperties; import org.springframework.cloud.stream.binder.rabbit.properties.RabbitExtendedBindingProperties; import org.springframework.cloud.stream.binder.rabbit.provisioning.RabbitExchangeQueueProvisioner; -import org.springframework.cloud.stream.config.codec.kryo.KryoCodecAutoConfiguration; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Import; -import org.springframework.integration.codec.Codec; /** * Configuration class for RabbitMQ message channel binder. * * @author David Turanski + * @author Vinicius Carvalho */ @Configuration -@Import({PropertyPlaceholderAutoConfiguration.class, KryoCodecAutoConfiguration.class}) +@Import({PropertyPlaceholderAutoConfiguration.class}) @EnableConfigurationProperties({RabbitBinderConfigurationProperties.class, RabbitExtendedBindingProperties.class}) public class RabbitMessageChannelBinderConfiguration { - @Autowired - private Codec codec; @Autowired private ConnectionFactory rabbitConnectionFactory; @@ -65,7 +62,6 @@ public class RabbitMessageChannelBinderConfiguration { RabbitMessageChannelBinder rabbitMessageChannelBinder() { RabbitMessageChannelBinder binder = new RabbitMessageChannelBinder(rabbitConnectionFactory, rabbitProperties, provisioningProvider()); - binder.setCodec(codec); binder.setAdminAddresses(rabbitBinderConfigurationProperties.getAdminAddresses()); binder.setCompressingPostProcessor(gZipPostProcessor()); binder.setDecompressingPostProcessor(deCompressingPostProcessor()); 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 31f3cb7bf..4b5854776 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 @@ -89,6 +89,7 @@ import org.springframework.messaging.SubscribableChannel; import org.springframework.messaging.support.ErrorMessage; import org.springframework.messaging.support.GenericMessage; import org.springframework.retry.support.RetryTemplate; +import org.springframework.util.MimeTypeUtils; import org.springframework.util.ReflectionUtils; import com.rabbitmq.http.client.domain.QueueInfo; @@ -159,22 +160,7 @@ public class RabbitBinderTests extends createProducerProperties()); Binding consumerBinding = binder.bindConsumer("bad.0", "test", moduleInputChannel, createConsumerProperties()); - - ConnectionFactory producerConnectionFactory = - TestUtils.getPropertyValue(producerBinding, "lifecycle.amqpTemplate.connectionFactory", - ConnectionFactory.class); - - ConnectionFactory consumerConnectionFactory = - TestUtils.getPropertyValue(consumerBinding, "lifecycle.messageListenerContainer.connectionFactory", - ConnectionFactory.class); - - assertThat(producerConnectionFactory).isNotSameAs(consumerConnectionFactory); - - assertThat(producerConnectionFactory.createConnection()) - .isNotEqualTo(consumerConnectionFactory.createConnection()); - - Message message = MessageBuilder.withPayload("bad").setHeader(MessageHeaders.CONTENT_TYPE, "foo/bar") - .build(); + Message message = MessageBuilder.withPayload("bad".getBytes()).setHeader(MessageHeaders.CONTENT_TYPE, "foo/bar").build(); final CountDownLatch latch = new CountDownLatch(3); moduleInputChannel.subscribe(new MessageHandler() { @@ -201,7 +187,7 @@ public class RabbitBinderTests extends ExtendedProducerProperties producerProps = createProducerProperties(); producerProps.setErrorChannelEnabled(true); Binding producerBinding = binder.bindProducer("ec.0", moduleOutputChannel, producerProps); - final Message message = MessageBuilder.withPayload("bad").setHeader(MessageHeaders.CONTENT_TYPE, "foo/bar") + final Message message = MessageBuilder.withPayload("bad".getBytes()).setHeader(MessageHeaders.CONTENT_TYPE, "foo/bar") .build(); SubscribableChannel ec = binder.getApplicationContext().getBean("ec.0.errors", SubscribableChannel.class); final AtomicReference> errorMessage = new AtomicReference<>(); @@ -1160,32 +1146,32 @@ public class RabbitBinderTests extends proxy.start(); - moduleOutputChannel.send(new GenericMessage<>("foo")); - Message message = moduleInputChannel.receive(20000); + moduleOutputChannel.send(MessageBuilder.withPayload("foo").setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.TEXT_PLAIN).build()); + Message message = moduleInputChannel.receive(10000); assertThat(message).isNotNull(); assertThat(message.getPayload()).isNotNull(); - noDLQOutputChannel.send(new GenericMessage<>("bar")); + noDLQOutputChannel.send(MessageBuilder.withPayload("bar").setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.TEXT_PLAIN).build()); message = noDLQInputChannel.receive(10000); assertThat(message); - assertThat(message.getPayload()).isEqualTo("bar"); + assertThat(message.getPayload()).isEqualTo("bar".getBytes()); - outputChannel.send(new GenericMessage<>("baz")); + outputChannel.send(MessageBuilder.withPayload("baz").setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.TEXT_PLAIN).build()); message = pubSubInputChannel.receive(10000); assertThat(message); - assertThat(message.getPayload()).isEqualTo("baz"); + assertThat(message.getPayload()).isEqualTo("baz".getBytes()); message = durablePubSubInputChannel.receive(10000); assertThat(message).isNotNull(); - assertThat(message.getPayload()).isEqualTo("baz"); + assertThat(message.getPayload()).isEqualTo("baz".getBytes()); - partOutputChannel.send(new GenericMessage<>("0")); - partOutputChannel.send(new GenericMessage<>("1")); + partOutputChannel.send(MessageBuilder.withPayload("0").setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.TEXT_PLAIN).build()); + partOutputChannel.send(MessageBuilder.withPayload("1").setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.TEXT_PLAIN).build()); message = partInputChannel0.receive(10000); assertThat(message).isNotNull(); - assertThat(message.getPayload()).isEqualTo("0"); + assertThat(message.getPayload()).isEqualTo("0".getBytes()); message = partInputChannel1.receive(10000); assertThat(message).isNotNull(); - assertThat(message.getPayload()).isEqualTo("1"); + assertThat(message.getPayload()).isEqualTo("1".getBytes()); late0ProducerBinding.unbind(); late0ConsumerBinding.unbind(); diff --git a/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/RabbitTestBinder.java b/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/RabbitTestBinder.java index d77fda360..e3e289711 100644 --- a/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/RabbitTestBinder.java +++ b/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/RabbitTestBinder.java @@ -32,7 +32,6 @@ import org.springframework.cloud.stream.binder.rabbit.properties.RabbitProducerP import org.springframework.cloud.stream.binder.rabbit.provisioning.RabbitExchangeQueueProvisioner; import org.springframework.context.annotation.AnnotationConfigApplicationContext; import org.springframework.context.annotation.Configuration; -import org.springframework.integration.codec.kryo.PojoCodec; import org.springframework.integration.config.EnableIntegration; import org.springframework.messaging.MessageChannel; @@ -64,7 +63,6 @@ public class RabbitTestBinder extends AbstractTestBinder