JMS Test Polishing

This commit is contained in:
Gary Russell
2015-11-26 13:29:24 -05:00
parent 4ff19def3f
commit 396b9146d2
4 changed files with 42 additions and 30 deletions

View File

@@ -1250,10 +1250,10 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp
}
else {
int n = 0;
while (this.replyDestination == null && n++ < 10) {
while (this.replyDestination == null && n++ < 100) {
logger.debug("Waiting for container to create destination");
try {
Thread.sleep(1000);
Thread.sleep(100);
}
catch (InterruptedException e) {
Thread.currentThread().interrupt();
@@ -1292,6 +1292,9 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp
@Override
protected void recoverAfterListenerSetupFailure() {
if (logger.isDebugEnabled()) {
logger.debug("recoverAfterListenerSetupFailure for dest: " + this.replyDestination);
}
this.replyDestination = null;
super.recoverAfterListenerSetupFailure();
}

View File

@@ -88,7 +88,7 @@ public class OutboundGatewayFunctionTests extends LogAdjustingTestSupport {
.thenReturn(scheduler);
final JmsOutboundGateway gateway = new JmsOutboundGateway();
gateway.setBeanFactory(beanFactory);
gateway.setConnectionFactory(getGatewayConnectionFactory());
gateway.setConnectionFactory(getConnectionFactory());
gateway.setRequestDestination(requestQueue1);
gateway.setReplyDestination(replyQueue1);
gateway.setCorrelationKey("JMSCorrelationID");
@@ -112,7 +112,7 @@ public class OutboundGatewayFunctionTests extends LogAdjustingTestSupport {
});
assertTrue(latch1.await(10, TimeUnit.SECONDS));
JmsTemplate template = new JmsTemplate();
template.setConnectionFactory(getTemplateConnectionFactory());
template.setConnectionFactory(getConnectionFactory());
template.setReceiveTimeout(5000);
javax.jms.Message request = template.receive(requestQueue1);
assertNotNull(request);
@@ -140,7 +140,7 @@ public class OutboundGatewayFunctionTests extends LogAdjustingTestSupport {
.thenReturn(scheduler);
final JmsOutboundGateway gateway = new JmsOutboundGateway();
gateway.setBeanFactory(beanFactory);
gateway.setConnectionFactory(getGatewayConnectionFactory());
gateway.setConnectionFactory(getConnectionFactory());
gateway.setRequestDestination(requestQueue2);
gateway.setReplyDestination(replyQueue2);
gateway.setUseReplyContainer(true);
@@ -163,7 +163,7 @@ public class OutboundGatewayFunctionTests extends LogAdjustingTestSupport {
});
assertTrue(latch1.await(10, TimeUnit.SECONDS));
JmsTemplate template = new JmsTemplate();
template.setConnectionFactory(getTemplateConnectionFactory());
template.setConnectionFactory(getConnectionFactory());
template.setReceiveTimeout(5000);
javax.jms.Message request = template.receive(requestQueue2);
assertNotNull(request);
@@ -192,7 +192,7 @@ public class OutboundGatewayFunctionTests extends LogAdjustingTestSupport {
.thenReturn(scheduler);
final JmsOutboundGateway gateway = new JmsOutboundGateway();
gateway.setBeanFactory(beanFactory);
gateway.setConnectionFactory(getGatewayConnectionFactory());
gateway.setConnectionFactory(getConnectionFactory());
gateway.setRequestDestination(requestQueue3);
gateway.setReplyDestinationName("reply3");
gateway.setCorrelationKey("JMSCorrelationID");
@@ -216,7 +216,7 @@ public class OutboundGatewayFunctionTests extends LogAdjustingTestSupport {
});
assertTrue(latch1.await(10, TimeUnit.SECONDS));
JmsTemplate template = new JmsTemplate();
template.setConnectionFactory(getTemplateConnectionFactory());
template.setConnectionFactory(getConnectionFactory());
template.setReceiveTimeout(5000);
javax.jms.Message request = template.receive(requestQueue3);
assertNotNull(request);
@@ -244,7 +244,7 @@ public class OutboundGatewayFunctionTests extends LogAdjustingTestSupport {
.thenReturn(scheduler);
final JmsOutboundGateway gateway = new JmsOutboundGateway();
gateway.setBeanFactory(beanFactory);
gateway.setConnectionFactory(getGatewayConnectionFactory());
gateway.setConnectionFactory(getConnectionFactory());
gateway.setRequestDestination(requestQueue4);
gateway.setReplyDestinationName("reply4");
gateway.setUseReplyContainer(true);
@@ -267,7 +267,7 @@ public class OutboundGatewayFunctionTests extends LogAdjustingTestSupport {
});
assertTrue(latch1.await(10, TimeUnit.SECONDS));
JmsTemplate template = new JmsTemplate();
template.setConnectionFactory(getTemplateConnectionFactory());
template.setConnectionFactory(getConnectionFactory());
template.setReceiveTimeout(5000);
javax.jms.Message request = template.receive(requestQueue4);
assertNotNull(request);
@@ -296,10 +296,11 @@ public class OutboundGatewayFunctionTests extends LogAdjustingTestSupport {
.thenReturn(scheduler);
final JmsOutboundGateway gateway = new JmsOutboundGateway();
gateway.setBeanFactory(beanFactory);
gateway.setConnectionFactory(getGatewayConnectionFactory());
gateway.setConnectionFactory(getConnectionFactory());
gateway.setRequestDestination(requestQueue5);
gateway.setCorrelationKey("JMSCorrelationID");
gateway.setUseReplyContainer(true);
gateway.setComponentName("testContainerWithTemporary.gateway");
gateway.afterPropertiesSet();
gateway.start();
final AtomicReference<Object> reply = new AtomicReference<Object>();
@@ -319,7 +320,7 @@ public class OutboundGatewayFunctionTests extends LogAdjustingTestSupport {
});
assertTrue(latch1.await(10, TimeUnit.SECONDS));
JmsTemplate template = new JmsTemplate();
template.setConnectionFactory(getTemplateConnectionFactory());
template.setConnectionFactory(getConnectionFactory());
template.setReceiveTimeout(5000);
javax.jms.Message request = template.receive(requestQueue5);
assertNotNull(request);
@@ -348,7 +349,7 @@ public class OutboundGatewayFunctionTests extends LogAdjustingTestSupport {
.thenReturn(scheduler);
final JmsOutboundGateway gateway = new JmsOutboundGateway();
gateway.setBeanFactory(beanFactory);
gateway.setConnectionFactory(getGatewayConnectionFactory());
gateway.setConnectionFactory(getConnectionFactory());
gateway.setRequestDestination(requestQueue6);
gateway.setUseReplyContainer(true);
gateway.afterPropertiesSet();
@@ -370,7 +371,7 @@ public class OutboundGatewayFunctionTests extends LogAdjustingTestSupport {
});
assertTrue(latch1.await(10, TimeUnit.SECONDS));
JmsTemplate template = new JmsTemplate();
template.setConnectionFactory(getTemplateConnectionFactory());
template.setConnectionFactory(getConnectionFactory());
template.setReceiveTimeout(5000);
javax.jms.Message request = template.receive(requestQueue6);
assertNotNull(request);
@@ -399,7 +400,7 @@ public class OutboundGatewayFunctionTests extends LogAdjustingTestSupport {
.thenReturn(scheduler);
final JmsOutboundGateway gateway = new JmsOutboundGateway();
gateway.setBeanFactory(beanFactory);
gateway.setConnectionFactory(getGatewayConnectionFactory());
gateway.setConnectionFactory(getConnectionFactory());
gateway.setRequestDestination(requestQueue7);
gateway.setReplyDestination(replyQueue7);
gateway.setCorrelationKey("JMSCorrelationID");
@@ -413,7 +414,7 @@ public class OutboundGatewayFunctionTests extends LogAdjustingTestSupport {
@Override
public void run() {
JmsTemplate template = new JmsTemplate();
template.setConnectionFactory(getTemplateConnectionFactory());
template.setConnectionFactory(getConnectionFactory());
template.setReceiveTimeout(20000);
receiveAndSend(template);
receiveAndSend(template);
@@ -453,14 +454,11 @@ public class OutboundGatewayFunctionTests extends LogAdjustingTestSupport {
assertFalse(container.isRunning());
}
private ConnectionFactory getTemplateConnectionFactory() {
ConnectionFactory amqConnectionFactory = new ActiveMQConnectionFactory("vm://localhost?broker.persistent=false");
return amqConnectionFactory;
}
private ConnectionFactory getGatewayConnectionFactory() {
ConnectionFactory amqConnectionFactory = new ActiveMQConnectionFactory("vm://localhost?broker.persistent=false");
return new CachingConnectionFactory(amqConnectionFactory);
private ConnectionFactory getConnectionFactory() {
ActiveMQConnectionFactory activeMQConnectionFactory = new ActiveMQConnectionFactory("vm://localhost?broker.persistent=false");
CachingConnectionFactory cachingConnectionFactory = new CachingConnectionFactory(activeMQConnectionFactory);
cachingConnectionFactory.setCacheConsumers(false);
return cachingConnectionFactory;
}
}

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2012 the original author or authors.
* Copyright 2002-2015 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.
@@ -21,17 +21,24 @@ import static org.junit.Assert.assertTrue;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.atomic.AtomicInteger;
import org.junit.Rule;
import org.junit.Test;
import org.springframework.context.support.ClassPathXmlApplicationContext;
import org.springframework.integration.gateway.RequestReplyExchanger;
import org.springframework.integration.jms.config.ActiveMqTestUtils;
import org.springframework.integration.test.support.LongRunningIntegrationTest;
import org.springframework.messaging.support.GenericMessage;
import org.springframework.util.StopWatch;
/**
* @author Oleg Zhurakousky
* @author Gary Russell
*/
public class MiscellaneousTests {
@Rule
public LongRunningIntegrationTest longRunning = new LongRunningIntegrationTest();
/**
* Asserts that receive-timeout is honored even if
* requests (once in process), takes less then receive-timeout value
@@ -51,18 +58,21 @@ public class MiscellaneousTests {
}
latch.await();
stopWatch.stop();
assertTrue(stopWatch.getTotalTimeMillis() <= 12000);
assertTrue(stopWatch.getTotalTimeMillis() <= 18000);
assertEquals(1, replies.get());
context.close();
}
private void exchange(final CountDownLatch latch, final RequestReplyExchanger gateway, final AtomicInteger replies) {
new Thread(new Runnable() {
@Override
public void run() {
try {
gateway.exchange(new GenericMessage<String>(""));
replies.incrementAndGet();
} catch (Exception e) {
}
catch (Exception e) {
//ignore
}
latch.countDown();

View File

@@ -7,18 +7,18 @@
http://www.springframework.org/schema/integration/jms http://www.springframework.org/schema/integration/jms/spring-integration-jms.xsd
http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans.xsd">
<int:gateway default-request-channel="in" default-request-timeout="10000" default-reply-timeout="10000"/>
<int:gateway default-request-channel="in" default-request-timeout="20000" default-reply-timeout="20000"/>
<int-jms:outbound-gateway request-channel="in"
connection-factory="connectionFactory"
request-destination-name="honorTimeoutQueue"
receive-timeout="10000"/>
receive-timeout="15000"/>
<int-jms:inbound-gateway request-channel="jmsIn"
request-destination-name="honorTimeoutQueue"
connection-factory="connectionFactory"
concurrent-consumers="1"
reply-timeout="10000"/>
reply-timeout="15000"/>
<int:chain input-channel="jmsIn">
<int:header-enricher>
@@ -35,6 +35,7 @@
<property name="brokerURL" value="vm://localhost?broker.persistent=false"/>
</bean>
</property>
<property name="cacheConsumers" value="false"/>
</bean>
</beans>