JmsInboundGateway now delegates to a MessageListener or SessionAwareMessageListener instance. The 'jms:inbound-gateway' parser now configures an instance of ChannelPublishingJmsMessageListener (part of INT-477 and INT-482).
This commit is contained in:
@@ -18,19 +18,16 @@ package org.springframework.integration.jms;
|
||||
|
||||
import javax.jms.ConnectionFactory;
|
||||
import javax.jms.Destination;
|
||||
import javax.jms.JMSException;
|
||||
import javax.jms.MessageProducer;
|
||||
import javax.jms.MessageListener;
|
||||
import javax.jms.Session;
|
||||
import javax.jms.Topic;
|
||||
|
||||
import org.springframework.beans.factory.DisposableBean;
|
||||
import org.springframework.core.task.TaskExecutor;
|
||||
import org.springframework.integration.core.Message;
|
||||
import org.springframework.integration.gateway.SimpleMessagingGateway;
|
||||
import org.springframework.integration.endpoint.AbstractEndpoint;
|
||||
import org.springframework.jms.listener.AbstractMessageListenerContainer;
|
||||
import org.springframework.jms.listener.DefaultMessageListenerContainer;
|
||||
import org.springframework.jms.listener.SessionAwareMessageListener;
|
||||
import org.springframework.jms.support.converter.MessageConverter;
|
||||
import org.springframework.transaction.PlatformTransactionManager;
|
||||
import org.springframework.util.Assert;
|
||||
|
||||
@@ -39,7 +36,9 @@ import org.springframework.util.Assert;
|
||||
*
|
||||
* @author Mark Fisher
|
||||
*/
|
||||
public class JmsInboundGateway extends SimpleMessagingGateway implements DisposableBean {
|
||||
public class JmsInboundGateway extends AbstractEndpoint implements DisposableBean {
|
||||
|
||||
private final Object listener;
|
||||
|
||||
private volatile AbstractMessageListenerContainer container;
|
||||
|
||||
@@ -51,14 +50,6 @@ public class JmsInboundGateway extends SimpleMessagingGateway implements Disposa
|
||||
|
||||
private volatile boolean pubSubDomain;
|
||||
|
||||
private volatile MessageConverter messageConverter;
|
||||
|
||||
private volatile JmsHeaderMapper headerMapper;
|
||||
|
||||
private volatile boolean extractRequestPayload = true;
|
||||
|
||||
private volatile boolean extractReplyPayload = true;
|
||||
|
||||
private volatile TaskExecutor taskExecutor;
|
||||
|
||||
private volatile PlatformTransactionManager transactionManager;
|
||||
@@ -76,6 +67,15 @@ public class JmsInboundGateway extends SimpleMessagingGateway implements Disposa
|
||||
private volatile int idleTaskExecutionLimit = 1;
|
||||
|
||||
|
||||
public JmsInboundGateway(Object listener) {
|
||||
Assert.notNull(listener, "listener must not be null");
|
||||
Assert.isTrue(listener instanceof MessageListener || listener instanceof SessionAwareMessageListener,
|
||||
"listener must implement either [" + MessageListener.class.getName()
|
||||
+ "] or [" + SessionAwareMessageListener.class.getName() + "]");
|
||||
this.listener = listener;
|
||||
}
|
||||
|
||||
|
||||
public void setContainer(AbstractMessageListenerContainer container) {
|
||||
this.container = container;
|
||||
}
|
||||
@@ -106,40 +106,6 @@ public class JmsInboundGateway extends SimpleMessagingGateway implements Disposa
|
||||
this.pubSubDomain = pubSubDomain;
|
||||
}
|
||||
|
||||
/**
|
||||
* Provide a {@link MessageConverter} implementation to use when
|
||||
* converting between JMS Messages and Spring Integration Messages.
|
||||
* If none is provided, a {@link HeaderMappingMessageConverter} will
|
||||
* be used and the {@link JmsHeaderMapper} instance provided to the
|
||||
* {@link #setHeaderMapper(JmsHeaderMapper)} method will be included
|
||||
* in the conversion process.
|
||||
*/
|
||||
public void setMessageConverter(MessageConverter messageConverter) {
|
||||
this.messageConverter = messageConverter;
|
||||
}
|
||||
|
||||
/**
|
||||
* Provide a {@link JmsHeaderMapper} implementation to use when
|
||||
* converting between JMS Messages and Spring Integration Messages.
|
||||
* If none is provided, a {@link DefaultJmsHeaderMapper} will be used.
|
||||
*
|
||||
* <p>This property will be ignored if a {@link MessageConverter} is
|
||||
* provided to the {@link #setMessageConverter(MessageConverter)} method.
|
||||
* However, you may provide your own implementation of the delegating
|
||||
* {@link HeaderMappingMessageConverter} implementation.
|
||||
*/
|
||||
public void setHeaderMapper(JmsHeaderMapper headerMapper) {
|
||||
this.headerMapper = headerMapper;
|
||||
}
|
||||
|
||||
public void setExtractRequestPayload(boolean extractRequestPayload) {
|
||||
this.extractRequestPayload = extractRequestPayload;
|
||||
}
|
||||
|
||||
public void setExtractReplyPayload(boolean extractReplyPayload) {
|
||||
this.extractReplyPayload = extractReplyPayload;
|
||||
}
|
||||
|
||||
public void setTaskExecutor(TaskExecutor taskExecutor) {
|
||||
this.taskExecutor = taskExecutor;
|
||||
}
|
||||
@@ -177,13 +143,7 @@ public class JmsInboundGateway extends SimpleMessagingGateway implements Disposa
|
||||
if (this.container == null) {
|
||||
this.container = createDefaultContainer();
|
||||
}
|
||||
if (this.messageConverter == null) {
|
||||
HeaderMappingMessageConverter hmmc = new HeaderMappingMessageConverter(null, this.headerMapper);
|
||||
hmmc.setExtractJmsMessageBody(this.extractRequestPayload);
|
||||
hmmc.setExtractIntegrationMessagePayload(this.extractReplyPayload);
|
||||
this.messageConverter = hmmc;
|
||||
}
|
||||
this.container.setMessageListener(new GatewayInvokingMessageListener());
|
||||
this.container.setMessageListener(this.listener);
|
||||
if (!this.container.isActive()) {
|
||||
this.container.afterPropertiesSet();
|
||||
}
|
||||
@@ -223,7 +183,6 @@ public class JmsInboundGateway extends SimpleMessagingGateway implements Disposa
|
||||
protected void doStart() {
|
||||
this.initialize();
|
||||
this.container.start();
|
||||
super.doStart();
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -231,7 +190,6 @@ public class JmsInboundGateway extends SimpleMessagingGateway implements Disposa
|
||||
if (this.container != null) {
|
||||
this.container.stop();
|
||||
}
|
||||
super.doStop();
|
||||
}
|
||||
|
||||
// DisposableBean implementation
|
||||
@@ -242,21 +200,4 @@ public class JmsInboundGateway extends SimpleMessagingGateway implements Disposa
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
private class GatewayInvokingMessageListener implements SessionAwareMessageListener {
|
||||
|
||||
public void onMessage(javax.jms.Message jmsMessage, Session session) throws JMSException {
|
||||
Object object = messageConverter.fromMessage(jmsMessage);
|
||||
Message<?> replyMessage = JmsInboundGateway.this.sendAndReceiveMessage(object);
|
||||
if (replyMessage != null) {
|
||||
javax.jms.Message jmsReply = messageConverter.toMessage(replyMessage, session);
|
||||
if (jmsReply.getJMSCorrelationID() == null) {
|
||||
jmsReply.setJMSCorrelationID(jmsMessage.getJMSMessageID());
|
||||
}
|
||||
MessageProducer producer = session.createProducer(jmsMessage.getJMSReplyTo());
|
||||
producer.send(jmsMessage.getJMSReplyTo(), jmsReply);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -22,9 +22,11 @@ import org.w3c.dom.Element;
|
||||
|
||||
import org.springframework.beans.factory.BeanCreationException;
|
||||
import org.springframework.beans.factory.support.BeanDefinitionBuilder;
|
||||
import org.springframework.beans.factory.support.BeanDefinitionReaderUtils;
|
||||
import org.springframework.beans.factory.xml.AbstractSingleBeanDefinitionParser;
|
||||
import org.springframework.beans.factory.xml.ParserContext;
|
||||
import org.springframework.integration.config.xml.IntegrationNamespaceUtils;
|
||||
import org.springframework.integration.jms.ChannelPublishingJmsMessageListener;
|
||||
import org.springframework.integration.jms.JmsInboundGateway;
|
||||
import org.springframework.util.StringUtils;
|
||||
|
||||
@@ -52,15 +54,10 @@ public class JmsInboundGatewayParser extends AbstractSingleBeanDefinitionParser
|
||||
|
||||
@Override
|
||||
protected void doParse(Element element, ParserContext parserContext, BeanDefinitionBuilder builder) {
|
||||
String listenerBeanName = this.parseMessageListener(element, parserContext);
|
||||
builder.addConstructorArgReference(listenerBeanName);
|
||||
String destination = element.getAttribute(JmsAdapterParserUtils.DESTINATION_ATTRIBUTE);
|
||||
String destinationName = element.getAttribute(JmsAdapterParserUtils.DESTINATION_NAME_ATTRIBUTE);
|
||||
String messageConverter = element.getAttribute(JmsAdapterParserUtils.MESSAGE_CONVERTER_ATTRIBUTE);
|
||||
if (StringUtils.hasText(element.getAttribute(JmsAdapterParserUtils.JMS_TEMPLATE_ATTRIBUTE))) {
|
||||
throw new BeanCreationException(JmsInboundGateway.class.getSimpleName() +
|
||||
" does not accept a '" + JmsAdapterParserUtils.JMS_TEMPLATE_ATTRIBUTE +
|
||||
"' reference. One of '" + JmsAdapterParserUtils.DESTINATION_ATTRIBUTE + "' or '" +
|
||||
JmsAdapterParserUtils.DESTINATION_NAME_ATTRIBUTE + "' must be provided.");
|
||||
}
|
||||
if (StringUtils.hasText(destination) || StringUtils.hasText(destinationName)) {
|
||||
builder.addPropertyReference(JmsAdapterParserUtils.CONNECTION_FACTORY_PROPERTY,
|
||||
JmsAdapterParserUtils.determineConnectionFactoryBeanName(element));
|
||||
@@ -75,9 +72,6 @@ public class JmsInboundGatewayParser extends AbstractSingleBeanDefinitionParser
|
||||
throw new BeanCreationException("One of '" + JmsAdapterParserUtils.DESTINATION_ATTRIBUTE +
|
||||
"' or '" + JmsAdapterParserUtils.DESTINATION_NAME_ATTRIBUTE + "' must be provided.");
|
||||
}
|
||||
if (StringUtils.hasText(messageConverter)) {
|
||||
builder.addPropertyReference(JmsAdapterParserUtils.MESSAGE_CONVERTER_PROPERTY, messageConverter);
|
||||
}
|
||||
Integer acknowledgeMode = JmsAdapterParserUtils.parseAcknowledgeMode(element);
|
||||
if (acknowledgeMode != null) {
|
||||
if (acknowledgeMode.intValue() == Session.SESSION_TRANSACTED) {
|
||||
@@ -87,15 +81,7 @@ public class JmsInboundGatewayParser extends AbstractSingleBeanDefinitionParser
|
||||
builder.addPropertyValue("sessionAcknowledgeMode", acknowledgeMode);
|
||||
}
|
||||
}
|
||||
IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "message-converter");
|
||||
IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "header-mapper");
|
||||
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "extract-request-payload");
|
||||
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "extract-reply-payload");
|
||||
IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "transaction-manager");
|
||||
IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "request-channel");
|
||||
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "request-timeout");
|
||||
IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "reply-channel");
|
||||
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "reply-timeout");
|
||||
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "pub-sub-domain");
|
||||
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "concurrent-consumers");
|
||||
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "max-concurrent-consumers");
|
||||
@@ -103,4 +89,17 @@ public class JmsInboundGatewayParser extends AbstractSingleBeanDefinitionParser
|
||||
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "idle-task-execution-limit");
|
||||
}
|
||||
|
||||
private String parseMessageListener(Element element, ParserContext parserContext) {
|
||||
BeanDefinitionBuilder builder = BeanDefinitionBuilder.genericBeanDefinition(ChannelPublishingJmsMessageListener.class);
|
||||
builder.addPropertyValue("expectReply", true);
|
||||
IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "request-channel");
|
||||
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "request-timeout");
|
||||
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "reply-timeout");
|
||||
IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "message-converter");
|
||||
IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "header-mapper");
|
||||
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "extract-request-payload");
|
||||
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "extract-reply-payload");
|
||||
return BeanDefinitionReaderUtils.registerWithGeneratedName(builder.getBeanDefinition(), parserContext.getRegistry());
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -85,6 +85,7 @@ public class JmsInboundGatewayParserTests {
|
||||
"jmsGatewaysWithExtractPayloadAttributes.xml", this.getClass());
|
||||
JmsInboundGateway gateway = (JmsInboundGateway) context.getBean("defaultGateway");
|
||||
DirectFieldAccessor accessor = new DirectFieldAccessor(gateway);
|
||||
accessor = new DirectFieldAccessor(accessor.getPropertyValue("listener"));
|
||||
assertEquals(Boolean.TRUE, accessor.getPropertyValue("extractReplyPayload"));
|
||||
}
|
||||
|
||||
@@ -94,6 +95,7 @@ public class JmsInboundGatewayParserTests {
|
||||
"jmsGatewaysWithExtractPayloadAttributes.xml", this.getClass());
|
||||
JmsInboundGateway gateway = (JmsInboundGateway) context.getBean("extractReplyPayloadTrue");
|
||||
DirectFieldAccessor accessor = new DirectFieldAccessor(gateway);
|
||||
accessor = new DirectFieldAccessor(accessor.getPropertyValue("listener"));
|
||||
assertEquals(Boolean.TRUE, accessor.getPropertyValue("extractReplyPayload"));
|
||||
}
|
||||
|
||||
@@ -103,6 +105,7 @@ public class JmsInboundGatewayParserTests {
|
||||
"jmsGatewaysWithExtractPayloadAttributes.xml", this.getClass());
|
||||
JmsInboundGateway gateway = (JmsInboundGateway) context.getBean("extractReplyPayloadFalse");
|
||||
DirectFieldAccessor accessor = new DirectFieldAccessor(gateway);
|
||||
accessor = new DirectFieldAccessor(accessor.getPropertyValue("listener"));
|
||||
assertEquals(Boolean.FALSE, accessor.getPropertyValue("extractReplyPayload"));
|
||||
}
|
||||
|
||||
@@ -112,6 +115,7 @@ public class JmsInboundGatewayParserTests {
|
||||
"jmsGatewaysWithExtractPayloadAttributes.xml", this.getClass());
|
||||
JmsInboundGateway gateway = (JmsInboundGateway) context.getBean("extractRequestPayloadTrue");
|
||||
DirectFieldAccessor accessor = new DirectFieldAccessor(gateway);
|
||||
accessor = new DirectFieldAccessor(accessor.getPropertyValue("listener"));
|
||||
assertEquals(Boolean.TRUE, accessor.getPropertyValue("extractRequestPayload"));
|
||||
}
|
||||
|
||||
@@ -121,6 +125,7 @@ public class JmsInboundGatewayParserTests {
|
||||
"jmsGatewaysWithExtractPayloadAttributes.xml", this.getClass());
|
||||
JmsInboundGateway gateway = (JmsInboundGateway) context.getBean("extractRequestPayloadFalse");
|
||||
DirectFieldAccessor accessor = new DirectFieldAccessor(gateway);
|
||||
accessor = new DirectFieldAccessor(accessor.getPropertyValue("listener"));
|
||||
assertEquals(Boolean.FALSE, accessor.getPropertyValue("extractRequestPayload"));
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user