INT-813, INT-965 refactoring JmsDestinationBackedMessageChannel into SubscribableJmsChannel and PollableJmsChannel with a common FactoryBean for creation

This commit is contained in:
Mark Fisher
2010-09-22 23:49:39 -04:00
parent d7e1041ed4
commit 55d9366c4b
10 changed files with 822 additions and 660 deletions

View File

@@ -36,16 +36,20 @@ import org.apache.activemq.command.ActiveMQTopic;
import org.junit.Before;
import org.junit.Test;
import org.springframework.beans.DirectFieldAccessor;
import org.springframework.beans.factory.support.BeanDefinitionBuilder;
import org.springframework.context.support.StaticApplicationContext;
import org.springframework.integration.Message;
import org.springframework.integration.core.MessageHandler;
import org.springframework.integration.jms.config.JmsChannelFactoryBean;
import org.springframework.integration.message.GenericMessage;
import org.springframework.jms.listener.AbstractMessageListenerContainer;
import org.springframework.jms.listener.DefaultMessageListenerContainer;
/**
* @author Mark Fisher
*/
public class JmsDestinationBackedMessageChannelTests {
public class SubscribableJmsChannelTests {
private static final int TIMEOUT = 30000;
@@ -82,8 +86,11 @@ public class JmsDestinationBackedMessageChannelTests {
latch.countDown();
}
};
JmsDestinationBackedMessageChannel channel =
new JmsDestinationBackedMessageChannel(this.connectionFactory, this.queue);
JmsChannelFactoryBean factoryBean = new JmsChannelFactoryBean(true);
factoryBean.setConnectionFactory(this.connectionFactory);
factoryBean.setDestination(this.queue);
factoryBean.afterPropertiesSet();
SubscribableJmsChannel channel = (SubscribableJmsChannel) factoryBean.getObject();
channel.afterPropertiesSet();
channel.start();
channel.subscribe(handler1);
@@ -117,13 +124,16 @@ public class JmsDestinationBackedMessageChannelTests {
latch.countDown();
}
};
JmsDestinationBackedMessageChannel channel =
new JmsDestinationBackedMessageChannel(this.connectionFactory, this.topic);
JmsChannelFactoryBean factoryBean = new JmsChannelFactoryBean(true);
factoryBean.setConnectionFactory(this.connectionFactory);
factoryBean.setDestination(this.topic);
factoryBean.afterPropertiesSet();
SubscribableJmsChannel channel = (SubscribableJmsChannel) factoryBean.getObject();
channel.afterPropertiesSet();
channel.subscribe(handler1);
channel.subscribe(handler2);
channel.start();
if (!channel.waitRegisteredWithDestination(10000)) {
if (!waitUntilRegisteredWithDestination(channel, 10000)) {
fail("Listener failed to subscribe to topic");
}
channel.send(new GenericMessage<String>("foo"));
@@ -155,8 +165,12 @@ public class JmsDestinationBackedMessageChannelTests {
latch.countDown();
}
};
JmsDestinationBackedMessageChannel channel =
new JmsDestinationBackedMessageChannel(this.connectionFactory, "dynamicQueue", false);
JmsChannelFactoryBean factoryBean = new JmsChannelFactoryBean(true);
factoryBean.setConnectionFactory(this.connectionFactory);
factoryBean.setDestinationName("dynamicQueue");
factoryBean.setPubSubDomain(false);
factoryBean.afterPropertiesSet();
SubscribableJmsChannel channel = (SubscribableJmsChannel) factoryBean.getObject();
channel.afterPropertiesSet();
channel.start();
channel.subscribe(handler1);
@@ -190,11 +204,16 @@ public class JmsDestinationBackedMessageChannelTests {
latch.countDown();
}
};
JmsDestinationBackedMessageChannel channel =
new JmsDestinationBackedMessageChannel(this.connectionFactory, "dynamicTopic", true);
JmsChannelFactoryBean factoryBean = new JmsChannelFactoryBean(true);
factoryBean.setConnectionFactory(this.connectionFactory);
factoryBean.setDestinationName("dynamicTopic");
factoryBean.setPubSubDomain(true);
factoryBean.afterPropertiesSet();
SubscribableJmsChannel channel = (SubscribableJmsChannel) factoryBean.getObject();
channel.afterPropertiesSet();
channel.start();
if (!channel.waitRegisteredWithDestination(10000)) {
if (!waitUntilRegisteredWithDestination(channel, 10000)) {
fail("Listener failed to subscribe to topic");
}
channel.subscribe(handler1);
@@ -211,15 +230,16 @@ public class JmsDestinationBackedMessageChannelTests {
channel.stop();
}
@Test
@Test //@Ignore
public void contextManagesLifecycle() {
BeanDefinitionBuilder builder = BeanDefinitionBuilder.genericBeanDefinition(JmsDestinationBackedMessageChannel.class);
builder.addConstructorArgValue(this.connectionFactory);
builder.addConstructorArgValue("dynamicQueue");
builder.addConstructorArgValue(false);
BeanDefinitionBuilder builder = BeanDefinitionBuilder.genericBeanDefinition(JmsChannelFactoryBean.class);
builder.addConstructorArgValue(true);
builder.addPropertyValue("connectionFactory", this.connectionFactory);
builder.addPropertyValue("destinationName", "dynamicQueue");
builder.addPropertyValue("pubSubDomain", false);
StaticApplicationContext context = new StaticApplicationContext();
context.registerBeanDefinition("channel", builder.getBeanDefinition());
JmsDestinationBackedMessageChannel channel = context.getBean("channel", JmsDestinationBackedMessageChannel.class);
SubscribableJmsChannel channel = context.getBean("channel", SubscribableJmsChannel.class);
assertFalse(channel.isRunning());
context.refresh();
assertTrue(channel.isRunning());
@@ -227,4 +247,38 @@ public class JmsDestinationBackedMessageChannelTests {
assertFalse(channel.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.
*/
private static boolean waitUntilRegisteredWithDestination(SubscribableJmsChannel channel, long timeout) {
AbstractMessageListenerContainer container =
(AbstractMessageListenerContainer) new DirectFieldAccessor(channel).getPropertyValue("container");
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;
}
}

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2009 the original author or authors.
* Copyright 2002-2010 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.
@@ -30,7 +30,7 @@ import org.junit.runner.RunWith;
import org.springframework.beans.DirectFieldAccessor;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.integration.MessageChannel;
import org.springframework.integration.jms.JmsDestinationBackedMessageChannel;
import org.springframework.integration.jms.SubscribableJmsChannel;
import org.springframework.jms.core.JmsTemplate;
import org.springframework.jms.listener.AbstractMessageListenerContainer;
import org.springframework.jms.listener.DefaultMessageListenerContainer;
@@ -75,8 +75,8 @@ public class JmsChannelParserTests {
@Test
public void queueReferenceChannel() {
assertEquals(JmsDestinationBackedMessageChannel.class, queueReferenceChannel.getClass());
JmsDestinationBackedMessageChannel channel = (JmsDestinationBackedMessageChannel) queueReferenceChannel;
assertEquals(SubscribableJmsChannel.class, queueReferenceChannel.getClass());
SubscribableJmsChannel channel = (SubscribableJmsChannel) queueReferenceChannel;
DirectFieldAccessor accessor = new DirectFieldAccessor(channel);
JmsTemplate jmsTemplate = (JmsTemplate) accessor.getPropertyValue("jmsTemplate");
AbstractMessageListenerContainer container = (AbstractMessageListenerContainer) accessor.getPropertyValue("container");
@@ -86,8 +86,8 @@ public class JmsChannelParserTests {
@Test
public void queueNameChannel() {
assertEquals(JmsDestinationBackedMessageChannel.class, queueNameChannel.getClass());
JmsDestinationBackedMessageChannel channel = (JmsDestinationBackedMessageChannel) queueNameChannel;
assertEquals(SubscribableJmsChannel.class, queueNameChannel.getClass());
SubscribableJmsChannel channel = (SubscribableJmsChannel) queueNameChannel;
DirectFieldAccessor accessor = new DirectFieldAccessor(channel);
JmsTemplate jmsTemplate = (JmsTemplate) accessor.getPropertyValue("jmsTemplate");
AbstractMessageListenerContainer container = (AbstractMessageListenerContainer) accessor.getPropertyValue("container");
@@ -97,8 +97,8 @@ public class JmsChannelParserTests {
@Test
public void queueNameWithResolverChannel() {
assertEquals(JmsDestinationBackedMessageChannel.class, queueNameWithResolverChannel.getClass());
JmsDestinationBackedMessageChannel channel = (JmsDestinationBackedMessageChannel) queueNameWithResolverChannel;
assertEquals(SubscribableJmsChannel.class, queueNameWithResolverChannel.getClass());
SubscribableJmsChannel channel = (SubscribableJmsChannel) queueNameWithResolverChannel;
DirectFieldAccessor accessor = new DirectFieldAccessor(channel);
JmsTemplate jmsTemplate = (JmsTemplate) accessor.getPropertyValue("jmsTemplate");
AbstractMessageListenerContainer container = (AbstractMessageListenerContainer) accessor.getPropertyValue("container");
@@ -108,8 +108,8 @@ public class JmsChannelParserTests {
@Test
public void topicReferenceChannel() {
assertEquals(JmsDestinationBackedMessageChannel.class, topicReferenceChannel.getClass());
JmsDestinationBackedMessageChannel channel = (JmsDestinationBackedMessageChannel) topicReferenceChannel;
assertEquals(SubscribableJmsChannel.class, topicReferenceChannel.getClass());
SubscribableJmsChannel channel = (SubscribableJmsChannel) topicReferenceChannel;
DirectFieldAccessor accessor = new DirectFieldAccessor(channel);
JmsTemplate jmsTemplate = (JmsTemplate) accessor.getPropertyValue("jmsTemplate");
AbstractMessageListenerContainer container = (AbstractMessageListenerContainer) accessor.getPropertyValue("container");
@@ -119,8 +119,8 @@ public class JmsChannelParserTests {
@Test
public void topicNameChannel() {
assertEquals(JmsDestinationBackedMessageChannel.class, topicNameChannel.getClass());
JmsDestinationBackedMessageChannel channel = (JmsDestinationBackedMessageChannel) topicNameChannel;
assertEquals(SubscribableJmsChannel.class, topicNameChannel.getClass());
SubscribableJmsChannel channel = (SubscribableJmsChannel) topicNameChannel;
DirectFieldAccessor accessor = new DirectFieldAccessor(channel);
JmsTemplate jmsTemplate = (JmsTemplate) accessor.getPropertyValue("jmsTemplate");
AbstractMessageListenerContainer container = (AbstractMessageListenerContainer) accessor.getPropertyValue("container");
@@ -130,8 +130,8 @@ public class JmsChannelParserTests {
@Test
public void topicNameWithResolverChannel() {
assertEquals(JmsDestinationBackedMessageChannel.class, topicNameWithResolverChannel.getClass());
JmsDestinationBackedMessageChannel channel = (JmsDestinationBackedMessageChannel) topicNameWithResolverChannel;
assertEquals(SubscribableJmsChannel.class, topicNameWithResolverChannel.getClass());
SubscribableJmsChannel channel = (SubscribableJmsChannel) topicNameWithResolverChannel;
DirectFieldAccessor accessor = new DirectFieldAccessor(channel);
JmsTemplate jmsTemplate = (JmsTemplate) accessor.getPropertyValue("jmsTemplate");
AbstractMessageListenerContainer container = (AbstractMessageListenerContainer) accessor.getPropertyValue("container");
@@ -141,8 +141,8 @@ public class JmsChannelParserTests {
@Test
public void channelWithConcurrencySettings() {
assertEquals(JmsDestinationBackedMessageChannel.class, channelWithConcurrencySettings.getClass());
JmsDestinationBackedMessageChannel channel = (JmsDestinationBackedMessageChannel) channelWithConcurrencySettings;
assertEquals(SubscribableJmsChannel.class, channelWithConcurrencySettings.getClass());
SubscribableJmsChannel channel = (SubscribableJmsChannel) channelWithConcurrencySettings;
DirectFieldAccessor accessor = new DirectFieldAccessor(channel);
DefaultMessageListenerContainer container = (DefaultMessageListenerContainer) accessor.getPropertyValue("container");
assertEquals(11, container.getConcurrentConsumers());