diff --git a/spring-integration-jms/src/main/java/org/springframework/integration/jms/MessageListenerContainerConfigurationSupport.java b/spring-integration-jms/src/main/java/org/springframework/integration/jms/MessageListenerContainerConfigurationSupport.java index 89d60fb612..f82721dad3 100644 --- a/spring-integration-jms/src/main/java/org/springframework/integration/jms/MessageListenerContainerConfigurationSupport.java +++ b/spring-integration-jms/src/main/java/org/springframework/integration/jms/MessageListenerContainerConfigurationSupport.java @@ -345,6 +345,38 @@ abstract class MessageListenerContainerConfigurationSupport implements Initializ return this.initialized && this.container.isRunning(); } + /** + * Blocks until the listener container has subscribed; if the container does not support + * this test, or the caching mode is incompatible, true is returned. Otherwise blocks + * until timeout milliseconds have passed, or the consumer has registered. + * @see DefaultMessageListenerContainer.isRegisteredWithDestination() + * @param timeout Timeout in milliseconds. + * @return True if a subscriber has connected or the container/attributes does not support + * the test. False if a valid container does not have a registered consumer within + * timeout milliseconds. + */ + public boolean waitRegisteredWithDestination(long timeout) { + AbstractMessageListenerContainer container = this.getListenerContainer(); + if (container instanceof DefaultMessageListenerContainer) { + DefaultMessageListenerContainer listenerContainer = + (DefaultMessageListenerContainer) container; + if (listenerContainer.getCacheLevel() != DefaultMessageListenerContainer.CACHE_CONSUMER) { + return true; + } + while (timeout > 0) { + if (listenerContainer.isRegisteredWithDestination()) { + return true; + } + try { + Thread.sleep(100); + } catch (InterruptedException e) { } + timeout -= 100; + } + return false; + } + return true; + } + public void start() { this.getListenerContainer().start(); } diff --git a/spring-integration-jms/src/test/java/org/springframework/integration/jms/JmsDestinationBackedMessageChannelTests.java b/spring-integration-jms/src/test/java/org/springframework/integration/jms/JmsDestinationBackedMessageChannelTests.java index e145269d54..748c52a814 100644 --- a/spring-integration-jms/src/test/java/org/springframework/integration/jms/JmsDestinationBackedMessageChannelTests.java +++ b/spring-integration-jms/src/test/java/org/springframework/integration/jms/JmsDestinationBackedMessageChannelTests.java @@ -117,7 +117,9 @@ public class JmsDestinationBackedMessageChannelTests { channel.subscribe(handler1); channel.subscribe(handler2); channel.start(); - Thread.sleep(5000); // allow time for listener to subscribe + if (!channel.waitRegisteredWithDestination(10000)) { + fail("Listener failed to subscribe to topic"); + } channel.send(new StringMessage("foo")); channel.send(new StringMessage("bar")); latch.await(TIMEOUT, TimeUnit.MILLISECONDS);