diff --git a/org.springframework.integration.jms/src/main/java/org/springframework/integration/jms/JmsInboundGateway.java b/org.springframework.integration.jms/src/main/java/org/springframework/integration/jms/JmsInboundGateway.java index e6ee78f6f7..542286800e 100644 --- a/org.springframework.integration.jms/src/main/java/org/springframework/integration/jms/JmsInboundGateway.java +++ b/org.springframework.integration.jms/src/main/java/org/springframework/integration/jms/JmsInboundGateway.java @@ -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. - * - *
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); - } - } - } - } diff --git a/org.springframework.integration.jms/src/main/java/org/springframework/integration/jms/config/JmsInboundGatewayParser.java b/org.springframework.integration.jms/src/main/java/org/springframework/integration/jms/config/JmsInboundGatewayParser.java index 90be775a28..272f5d768d 100644 --- a/org.springframework.integration.jms/src/main/java/org/springframework/integration/jms/config/JmsInboundGatewayParser.java +++ b/org.springframework.integration.jms/src/main/java/org/springframework/integration/jms/config/JmsInboundGatewayParser.java @@ -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()); + } + } diff --git a/org.springframework.integration.jms/src/test/java/org/springframework/integration/jms/config/JmsInboundGatewayParserTests.java b/org.springframework.integration.jms/src/test/java/org/springframework/integration/jms/config/JmsInboundGatewayParserTests.java index 42d508f08c..ca4a52c432 100644 --- a/org.springframework.integration.jms/src/test/java/org/springframework/integration/jms/config/JmsInboundGatewayParserTests.java +++ b/org.springframework.integration.jms/src/test/java/org/springframework/integration/jms/config/JmsInboundGatewayParserTests.java @@ -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")); }