|
|
|
|
@@ -44,9 +44,9 @@ import com.rabbitmq.http.client.domain.BindingInfo;
|
|
|
|
|
import com.rabbitmq.http.client.domain.ExchangeInfo;
|
|
|
|
|
import com.rabbitmq.http.client.domain.QueueInfo;
|
|
|
|
|
import org.apache.commons.logging.Log;
|
|
|
|
|
import org.junit.Rule;
|
|
|
|
|
import org.junit.Test;
|
|
|
|
|
import org.junit.rules.TestName;
|
|
|
|
|
import org.junit.jupiter.api.Test;
|
|
|
|
|
import org.junit.jupiter.api.TestInfo;
|
|
|
|
|
import org.junit.jupiter.api.extension.RegisterExtension;
|
|
|
|
|
import org.mockito.ArgumentCaptor;
|
|
|
|
|
|
|
|
|
|
import org.springframework.amqp.AmqpIOException;
|
|
|
|
|
@@ -160,11 +160,9 @@ public class RabbitBinderTests extends
|
|
|
|
|
|
|
|
|
|
private int maxStackTraceSize;
|
|
|
|
|
|
|
|
|
|
@Rule
|
|
|
|
|
public RabbitTestSupport rabbitAvailableRule = new RabbitTestSupport(true);
|
|
|
|
|
|
|
|
|
|
@Rule
|
|
|
|
|
public TestName testName = new TestName();
|
|
|
|
|
@RegisterExtension
|
|
|
|
|
RabbitTestSupport rabbitAvailableRule = new RabbitTestSupport(true);
|
|
|
|
|
|
|
|
|
|
@Override
|
|
|
|
|
protected RabbitTestBinder getBinder() {
|
|
|
|
|
@@ -181,10 +179,11 @@ public class RabbitBinderTests extends
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@Override
|
|
|
|
|
protected ExtendedProducerProperties<RabbitProducerProperties> createProducerProperties() {
|
|
|
|
|
protected ExtendedProducerProperties<RabbitProducerProperties> createProducerProperties(TestInfo testInfo) {
|
|
|
|
|
ExtendedProducerProperties<RabbitProducerProperties> props = new ExtendedProducerProperties<>(
|
|
|
|
|
new RabbitProducerProperties());
|
|
|
|
|
if (testName.getMethodName().equals("testPartitionedModuleSpEL")) {
|
|
|
|
|
|
|
|
|
|
if (testInfo.getTestMethod().get().getName().equals("testPartitionedModuleSpEL")) {
|
|
|
|
|
props.getExtension().setRoutingKeyExpression(
|
|
|
|
|
spelExpressionParser.parseExpression("'part.0'"));
|
|
|
|
|
}
|
|
|
|
|
@@ -197,7 +196,7 @@ public class RabbitBinderTests extends
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@Test
|
|
|
|
|
public void testSendAndReceiveBad() throws Exception {
|
|
|
|
|
public void testSendAndReceiveBad(TestInfo testInfo) throws Exception {
|
|
|
|
|
RabbitTestBinder binder = getBinder();
|
|
|
|
|
final AtomicReference<AsyncConsumerStartedEvent> event = new AtomicReference<>();
|
|
|
|
|
binder.getApplicationContext().addApplicationListener(
|
|
|
|
|
@@ -207,7 +206,7 @@ public class RabbitBinderTests extends
|
|
|
|
|
DirectChannel moduleInputChannel = createBindableChannel("input",
|
|
|
|
|
new BindingProperties());
|
|
|
|
|
Binding<MessageChannel> producerBinding = binder.bindProducer("bad.0",
|
|
|
|
|
moduleOutputChannel, createProducerProperties());
|
|
|
|
|
moduleOutputChannel, createProducerProperties(testInfo));
|
|
|
|
|
assertThat(TestUtils.getPropertyValue(producerBinding,
|
|
|
|
|
"lifecycle.headersMappedLast", Boolean.class)).isTrue();
|
|
|
|
|
assertThat(
|
|
|
|
|
@@ -244,7 +243,7 @@ public class RabbitBinderTests extends
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@Test
|
|
|
|
|
public void testProducerErrorChannel() throws Exception {
|
|
|
|
|
public void testProducerErrorChannel(TestInfo testInfo) throws Exception {
|
|
|
|
|
RabbitTestBinder binder = getBinder();
|
|
|
|
|
CachingConnectionFactory ccf = this.rabbitAvailableRule.getResource();
|
|
|
|
|
ccf.setPublisherReturns(true);
|
|
|
|
|
@@ -252,7 +251,7 @@ public class RabbitBinderTests extends
|
|
|
|
|
ccf.resetConnection();
|
|
|
|
|
DirectChannel moduleOutputChannel = createBindableChannel("output",
|
|
|
|
|
new BindingProperties());
|
|
|
|
|
ExtendedProducerProperties<RabbitProducerProperties> producerProps = createProducerProperties();
|
|
|
|
|
ExtendedProducerProperties<RabbitProducerProperties> producerProps = createProducerProperties(testInfo);
|
|
|
|
|
producerProps.setErrorChannelEnabled(true);
|
|
|
|
|
Binding<MessageChannel> producerBinding = binder.bindProducer("ec.0",
|
|
|
|
|
moduleOutputChannel, producerProps);
|
|
|
|
|
@@ -325,7 +324,7 @@ public class RabbitBinderTests extends
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@Test
|
|
|
|
|
public void testProducerAckChannel() throws Exception {
|
|
|
|
|
public void testProducerAckChannel(TestInfo testInfo) throws Exception {
|
|
|
|
|
RabbitTestBinder binder = getBinder();
|
|
|
|
|
CachingConnectionFactory ccf = this.rabbitAvailableRule.getResource();
|
|
|
|
|
ccf.setPublisherReturns(true);
|
|
|
|
|
@@ -333,7 +332,7 @@ public class RabbitBinderTests extends
|
|
|
|
|
ccf.resetConnection();
|
|
|
|
|
DirectChannel moduleOutputChannel = createBindableChannel("output",
|
|
|
|
|
new BindingProperties());
|
|
|
|
|
ExtendedProducerProperties<RabbitProducerProperties> producerProps = createProducerProperties();
|
|
|
|
|
ExtendedProducerProperties<RabbitProducerProperties> producerProps = createProducerProperties(testInfo);
|
|
|
|
|
producerProps.setErrorChannelEnabled(true);
|
|
|
|
|
producerProps.getExtension().setConfirmAckChannel("acksChannel");
|
|
|
|
|
Binding<MessageChannel> producerBinding = binder.bindProducer("acks.0",
|
|
|
|
|
@@ -354,7 +353,7 @@ public class RabbitBinderTests extends
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@Test
|
|
|
|
|
public void testProducerConfirmHeader() throws Exception {
|
|
|
|
|
public void testProducerConfirmHeader(TestInfo testInfo) throws Exception {
|
|
|
|
|
RabbitTestBinder binder = getBinder();
|
|
|
|
|
CachingConnectionFactory ccf = this.rabbitAvailableRule.getResource();
|
|
|
|
|
ccf.setPublisherReturns(true);
|
|
|
|
|
@@ -362,7 +361,7 @@ public class RabbitBinderTests extends
|
|
|
|
|
ccf.resetConnection();
|
|
|
|
|
DirectChannel moduleOutputChannel = createBindableChannel("output",
|
|
|
|
|
new BindingProperties());
|
|
|
|
|
ExtendedProducerProperties<RabbitProducerProperties> producerProps = createProducerProperties();
|
|
|
|
|
ExtendedProducerProperties<RabbitProducerProperties> producerProps = createProducerProperties(testInfo);
|
|
|
|
|
producerProps.getExtension().setUseConfirmHeader(true);
|
|
|
|
|
Binding<MessageChannel> producerBinding = binder.bindProducer("confirms.0",
|
|
|
|
|
moduleOutputChannel, producerProps);
|
|
|
|
|
@@ -791,11 +790,11 @@ public class RabbitBinderTests extends
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@Test
|
|
|
|
|
public void testProducerProperties() throws Exception {
|
|
|
|
|
public void testProducerProperties(TestInfo testInfo) throws Exception {
|
|
|
|
|
RabbitTestBinder binder = getBinder();
|
|
|
|
|
Binding<MessageChannel> producerBinding = binder.bindProducer("props.0",
|
|
|
|
|
createBindableChannel("input", new BindingProperties()),
|
|
|
|
|
createProducerProperties());
|
|
|
|
|
createProducerProperties(testInfo));
|
|
|
|
|
Lifecycle endpoint = extractEndpoint(producerBinding);
|
|
|
|
|
MessageDeliveryMode mode = TestUtils.getPropertyValue(endpoint,
|
|
|
|
|
"defaultDeliveryMode", MessageDeliveryMode.class);
|
|
|
|
|
@@ -808,7 +807,7 @@ public class RabbitBinderTests extends
|
|
|
|
|
assertThat(TestUtils.getPropertyValue(endpoint, "amqpTemplate.transactional",
|
|
|
|
|
Boolean.class)).isFalse();
|
|
|
|
|
|
|
|
|
|
ExtendedProducerProperties<RabbitProducerProperties> producerProperties = createProducerProperties();
|
|
|
|
|
ExtendedProducerProperties<RabbitProducerProperties> producerProperties = createProducerProperties(testInfo);
|
|
|
|
|
this.applicationContext.registerBean("pkExtractor",
|
|
|
|
|
TestPartitionKeyExtractorClass.class, () -> new TestPartitionKeyExtractorClass());
|
|
|
|
|
this.applicationContext.registerBean("pkSelector",
|
|
|
|
|
@@ -1090,7 +1089,7 @@ public class RabbitBinderTests extends
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@Test
|
|
|
|
|
public void testAutoBindDLQPartionedConsumerFirst() throws Exception {
|
|
|
|
|
public void testAutoBindDLQPartionedConsumerFirst(TestInfo testInfo) throws Exception {
|
|
|
|
|
RabbitTestBinder binder = getBinder();
|
|
|
|
|
ExtendedConsumerProperties<RabbitConsumerProperties> properties = createConsumerProperties();
|
|
|
|
|
properties.getExtension().setPrefix("bindertest.");
|
|
|
|
|
@@ -1114,7 +1113,7 @@ public class RabbitBinderTests extends
|
|
|
|
|
Binding<MessageChannel> defaultConsumerBinding2 = binder.bindConsumer("partDLQ.0",
|
|
|
|
|
"default", new QueueChannel(), properties);
|
|
|
|
|
|
|
|
|
|
ExtendedProducerProperties<RabbitProducerProperties> producerProperties = createProducerProperties();
|
|
|
|
|
ExtendedProducerProperties<RabbitProducerProperties> producerProperties = createProducerProperties(testInfo);
|
|
|
|
|
producerProperties.getExtension().setPrefix("bindertest.");
|
|
|
|
|
this.applicationContext.registerBean("pkExtractor", PartitionTestSupport.class, () -> new PartitionTestSupport());
|
|
|
|
|
this.applicationContext.registerBean("pkSelector", PartitionTestSupport.class, () -> new PartitionTestSupport());
|
|
|
|
|
@@ -1192,19 +1191,19 @@ public class RabbitBinderTests extends
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@Test
|
|
|
|
|
public void testAutoBindDLQPartitionedConsumerFirstWithRepublishNoRetry()
|
|
|
|
|
public void testAutoBindDLQPartitionedConsumerFirstWithRepublishNoRetry(TestInfo testInfo)
|
|
|
|
|
throws Exception {
|
|
|
|
|
testAutoBindDLQPartionedConsumerFirstWithRepublishGuts(false);
|
|
|
|
|
testAutoBindDLQPartionedConsumerFirstWithRepublishGuts(false, testInfo);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@Test
|
|
|
|
|
public void testAutoBindDLQPartitionedConsumerFirstWithRepublishWithRetry()
|
|
|
|
|
public void testAutoBindDLQPartitionedConsumerFirstWithRepublishWithRetry(TestInfo testInfo)
|
|
|
|
|
throws Exception {
|
|
|
|
|
testAutoBindDLQPartionedConsumerFirstWithRepublishGuts(true);
|
|
|
|
|
testAutoBindDLQPartionedConsumerFirstWithRepublishGuts(true, testInfo);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private void testAutoBindDLQPartionedConsumerFirstWithRepublishGuts(
|
|
|
|
|
final boolean withRetry) throws Exception {
|
|
|
|
|
final boolean withRetry, TestInfo testInfo) throws Exception {
|
|
|
|
|
RabbitTestBinder binder = getBinder();
|
|
|
|
|
ExtendedConsumerProperties<RabbitConsumerProperties> properties = createConsumerProperties();
|
|
|
|
|
properties.getExtension().setPrefix("bindertest.");
|
|
|
|
|
@@ -1231,7 +1230,7 @@ public class RabbitBinderTests extends
|
|
|
|
|
Binding<MessageChannel> defaultConsumerBinding2 = binder
|
|
|
|
|
.bindConsumer("partPubDLQ.0", "default", new QueueChannel(), properties);
|
|
|
|
|
|
|
|
|
|
ExtendedProducerProperties<RabbitProducerProperties> producerProperties = createProducerProperties();
|
|
|
|
|
ExtendedProducerProperties<RabbitProducerProperties> producerProperties = createProducerProperties(testInfo);
|
|
|
|
|
producerProperties.getExtension().setPrefix("bindertest.");
|
|
|
|
|
producerProperties.getExtension().setAutoBindDlq(true);
|
|
|
|
|
this.applicationContext.registerBean("pkExtractor", PartitionTestSupport.class, () -> new PartitionTestSupport());
|
|
|
|
|
@@ -1349,9 +1348,9 @@ public class RabbitBinderTests extends
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@Test
|
|
|
|
|
public void testAutoBindDLQPartitionedProducerFirst() throws Exception {
|
|
|
|
|
public void testAutoBindDLQPartitionedProducerFirst(TestInfo testInfo) throws Exception {
|
|
|
|
|
RabbitTestBinder binder = getBinder();
|
|
|
|
|
ExtendedProducerProperties<RabbitProducerProperties> properties = createProducerProperties();
|
|
|
|
|
ExtendedProducerProperties<RabbitProducerProperties> properties = createProducerProperties(testInfo);
|
|
|
|
|
|
|
|
|
|
properties.getExtension().setPrefix("bindertest.");
|
|
|
|
|
properties.getExtension().setAutoBindDlq(true);
|
|
|
|
|
@@ -1664,9 +1663,9 @@ public class RabbitBinderTests extends
|
|
|
|
|
|
|
|
|
|
@SuppressWarnings("unchecked")
|
|
|
|
|
@Test
|
|
|
|
|
public void testBatchingAndCompression() throws Exception {
|
|
|
|
|
public void testBatchingAndCompression(TestInfo testInfo) throws Exception {
|
|
|
|
|
RabbitTestBinder binder = getBinder();
|
|
|
|
|
ExtendedProducerProperties<RabbitProducerProperties> producerProperties = createProducerProperties();
|
|
|
|
|
ExtendedProducerProperties<RabbitProducerProperties> producerProperties = createProducerProperties(testInfo);
|
|
|
|
|
producerProperties.getExtension()
|
|
|
|
|
.setDeliveryMode(MessageDeliveryMode.NON_PERSISTENT);
|
|
|
|
|
producerProperties.getExtension().setBatchingEnabled(true);
|
|
|
|
|
@@ -1726,9 +1725,9 @@ public class RabbitBinderTests extends
|
|
|
|
|
|
|
|
|
|
@SuppressWarnings("unchecked")
|
|
|
|
|
@Test
|
|
|
|
|
public void testProducerBatching() throws Exception {
|
|
|
|
|
public void testProducerBatching(TestInfo testInfo) throws Exception {
|
|
|
|
|
RabbitTestBinder binder = getBinder();
|
|
|
|
|
ExtendedProducerProperties<RabbitProducerProperties> producerProperties = createProducerProperties();
|
|
|
|
|
ExtendedProducerProperties<RabbitProducerProperties> producerProperties = createProducerProperties(testInfo);
|
|
|
|
|
producerProperties.getExtension()
|
|
|
|
|
.setDeliveryMode(MessageDeliveryMode.NON_PERSISTENT);
|
|
|
|
|
producerProperties.getExtension().setBatchingEnabled(true);
|
|
|
|
|
@@ -1768,9 +1767,9 @@ public class RabbitBinderTests extends
|
|
|
|
|
|
|
|
|
|
@SuppressWarnings("unchecked")
|
|
|
|
|
@Test
|
|
|
|
|
public void testConsumerBatching() throws Exception {
|
|
|
|
|
public void testConsumerBatching(TestInfo testInfo) throws Exception {
|
|
|
|
|
RabbitTestBinder binder = getBinder();
|
|
|
|
|
ExtendedProducerProperties<RabbitProducerProperties> producerProperties = createProducerProperties();
|
|
|
|
|
ExtendedProducerProperties<RabbitProducerProperties> producerProperties = createProducerProperties(testInfo);
|
|
|
|
|
producerProperties.getExtension()
|
|
|
|
|
.setDeliveryMode(MessageDeliveryMode.NON_PERSISTENT);
|
|
|
|
|
|
|
|
|
|
@@ -1806,9 +1805,9 @@ public class RabbitBinderTests extends
|
|
|
|
|
|
|
|
|
|
@SuppressWarnings("unchecked")
|
|
|
|
|
@Test
|
|
|
|
|
public void testInternalHeadersNotPropagated() throws Exception {
|
|
|
|
|
public void testInternalHeadersNotPropagated(TestInfo testInfo) throws Exception {
|
|
|
|
|
RabbitTestBinder binder = getBinder();
|
|
|
|
|
ExtendedProducerProperties<RabbitProducerProperties> producerProperties = createProducerProperties();
|
|
|
|
|
ExtendedProducerProperties<RabbitProducerProperties> producerProperties = createProducerProperties(testInfo);
|
|
|
|
|
producerProperties.getExtension()
|
|
|
|
|
.setDeliveryMode(MessageDeliveryMode.NON_PERSISTENT);
|
|
|
|
|
|
|
|
|
|
@@ -1848,7 +1847,7 @@ public class RabbitBinderTests extends
|
|
|
|
|
* queues.
|
|
|
|
|
*/
|
|
|
|
|
@Test
|
|
|
|
|
public void testLateBinding() throws Exception {
|
|
|
|
|
public void testLateBinding(TestInfo testInfo) throws Exception {
|
|
|
|
|
RabbitTestSupport.RabbitProxy proxy = new RabbitTestSupport.RabbitProxy();
|
|
|
|
|
CachingConnectionFactory cf = new CachingConnectionFactory("localhost",
|
|
|
|
|
proxy.getPort());
|
|
|
|
|
@@ -1857,7 +1856,7 @@ public class RabbitBinderTests extends
|
|
|
|
|
new RabbitProperties(), new RabbitExchangeQueueProvisioner(cf));
|
|
|
|
|
RabbitTestBinder binder = new RabbitTestBinder(cf, rabbitBinder);
|
|
|
|
|
|
|
|
|
|
ExtendedProducerProperties<RabbitProducerProperties> producerProperties = createProducerProperties();
|
|
|
|
|
ExtendedProducerProperties<RabbitProducerProperties> producerProperties = createProducerProperties(testInfo);
|
|
|
|
|
producerProperties.getExtension().setPrefix("latebinder.");
|
|
|
|
|
producerProperties.getExtension().setAutoBindDlq(true);
|
|
|
|
|
producerProperties.getExtension().setTransacted(true);
|
|
|
|
|
@@ -1897,7 +1896,7 @@ public class RabbitBinderTests extends
|
|
|
|
|
Binding<MessageChannel> partlate0Consumer1Binding = binder.bindConsumer(
|
|
|
|
|
"partlate.0", "test", partInputChannel1, partLateConsumerProperties);
|
|
|
|
|
|
|
|
|
|
ExtendedProducerProperties<RabbitProducerProperties> noDlqProducerProperties = createProducerProperties();
|
|
|
|
|
ExtendedProducerProperties<RabbitProducerProperties> noDlqProducerProperties = createProducerProperties(testInfo);
|
|
|
|
|
noDlqProducerProperties.getExtension().setPrefix("latebinder.");
|
|
|
|
|
MessageChannel noDLQOutputChannel = createBindableChannel("output",
|
|
|
|
|
createProducerBindingProperties(noDlqProducerProperties));
|
|
|
|
|
@@ -2021,9 +2020,9 @@ public class RabbitBinderTests extends
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@Test
|
|
|
|
|
public void testRoutingKeyExpression() throws Exception {
|
|
|
|
|
public void testRoutingKeyExpression(TestInfo testInfo) throws Exception {
|
|
|
|
|
RabbitTestBinder binder = getBinder();
|
|
|
|
|
ExtendedProducerProperties<RabbitProducerProperties> producerProperties = createProducerProperties();
|
|
|
|
|
ExtendedProducerProperties<RabbitProducerProperties> producerProperties = createProducerProperties(testInfo);
|
|
|
|
|
producerProperties.getExtension().setRoutingKeyExpression(
|
|
|
|
|
spelExpressionParser.parseExpression("payload.field"));
|
|
|
|
|
|
|
|
|
|
@@ -2064,9 +2063,9 @@ public class RabbitBinderTests extends
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@Test
|
|
|
|
|
public void testRoutingKeyExpressionPartitionedAndDelay() throws Exception {
|
|
|
|
|
public void testRoutingKeyExpressionPartitionedAndDelay(TestInfo testInfo) throws Exception {
|
|
|
|
|
RabbitTestBinder binder = getBinder();
|
|
|
|
|
ExtendedProducerProperties<RabbitProducerProperties> producerProperties = createProducerProperties();
|
|
|
|
|
ExtendedProducerProperties<RabbitProducerProperties> producerProperties = createProducerProperties(testInfo);
|
|
|
|
|
producerProperties.getExtension().setRoutingKeyExpression(
|
|
|
|
|
spelExpressionParser.parseExpression("#root.getPayload().field"));
|
|
|
|
|
// requires delayed message exchange plugin; tested locally
|
|
|
|
|
@@ -2265,9 +2264,9 @@ public class RabbitBinderTests extends
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@Test
|
|
|
|
|
public void testCustomBatchingStrategy() throws Exception {
|
|
|
|
|
public void testCustomBatchingStrategy(TestInfo testInfo) throws Exception {
|
|
|
|
|
RabbitTestBinder binder = getBinder();
|
|
|
|
|
ExtendedProducerProperties<RabbitProducerProperties> producerProperties = createProducerProperties();
|
|
|
|
|
ExtendedProducerProperties<RabbitProducerProperties> producerProperties = createProducerProperties(testInfo);
|
|
|
|
|
producerProperties.getExtension().setDeliveryMode(MessageDeliveryMode.NON_PERSISTENT);
|
|
|
|
|
producerProperties.getExtension().setBatchingEnabled(true);
|
|
|
|
|
producerProperties.getExtension().setBatchingStrategyBeanName("testCustomBatchingStrategy");
|
|
|
|
|
|