ContentType revamp
- Making tests compatible with new contentType behavior
This commit is contained in:
committed by
Soby Chacko
parent
d516cfd3db
commit
f4f1785e20
5
pom.xml
5
pom.xml
@@ -23,11 +23,6 @@
|
||||
<artifactId>spring-cloud-stream</artifactId>
|
||||
<version>${spring-cloud-stream.version}</version>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>org.springframework.cloud</groupId>
|
||||
<artifactId>spring-cloud-stream-codec</artifactId>
|
||||
<version>${spring-cloud-stream.version}</version>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>org.springframework.cloud</groupId>
|
||||
<artifactId>spring-cloud-stream-binder-rabbit</artifactId>
|
||||
|
||||
@@ -28,10 +28,6 @@
|
||||
<groupId>org.springframework.cloud</groupId>
|
||||
<artifactId>spring-cloud-stream</artifactId>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>org.springframework.cloud</groupId>
|
||||
<artifactId>spring-cloud-stream-codec</artifactId>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>org.springframework.boot</groupId>
|
||||
<artifactId>spring-boot-autoconfigure</artifactId>
|
||||
|
||||
@@ -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());
|
||||
|
||||
@@ -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<MessageChannel> 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<RabbitProducerProperties> producerProps = createProducerProperties();
|
||||
producerProps.setErrorChannelEnabled(true);
|
||||
Binding<MessageChannel> 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<Message<?>> 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();
|
||||
|
||||
@@ -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<RabbitMessageChannelBin
|
||||
public RabbitTestBinder(ConnectionFactory connectionFactory, RabbitMessageChannelBinder binder) {
|
||||
this.applicationContext = new AnnotationConfigApplicationContext(Config.class);
|
||||
binder.setApplicationContext(this.applicationContext);
|
||||
binder.setCodec(new PojoCodec());
|
||||
this.setBinder(binder);
|
||||
this.rabbitAdmin = new RabbitAdmin(connectionFactory);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user