diff --git a/org.springframework.integration.jms/src/main/java/org/springframework/integration/jms/ChannelPublishingJmsMessageListener.java b/org.springframework.integration.jms/src/main/java/org/springframework/integration/jms/ChannelPublishingJmsMessageListener.java index d1c9ccc1c4..5151bd73a4 100644 --- a/org.springframework.integration.jms/src/main/java/org/springframework/integration/jms/ChannelPublishingJmsMessageListener.java +++ b/org.springframework.integration.jms/src/main/java/org/springframework/integration/jms/ChannelPublishingJmsMessageListener.java @@ -62,6 +62,8 @@ public class ChannelPublishingJmsMessageListener implements SessionAwareMessageL private volatile int replyDeliveryMode = javax.jms.Message.DEFAULT_DELIVERY_MODE; + private volatile boolean explicitQosEnabledForReplies; + private volatile DestinationResolver destinationResolver = new DynamicDestinationResolver(); private volatile JmsHeaderMapper headerMapper; @@ -158,6 +160,14 @@ public class ChannelPublishingJmsMessageListener implements SessionAwareMessageL this.replyDeliveryMode = replyDeliveryPersistent ? DeliveryMode.PERSISTENT : DeliveryMode.NON_PERSISTENT; } + /** + * Specify whether explicit QoS should be enabled for replies + * (for timeToLive, priority, and deliveryMode settings). + */ + public void setExplicitQosEnabledForReplies(boolean explicitQosEnabledForReplies) { + this.explicitQosEnabledForReplies = explicitQosEnabledForReplies; + } + /** * Set the DestinationResolver that should be used to resolve reply * destination names for this listener. @@ -245,11 +255,14 @@ public class ChannelPublishingJmsMessageListener implements SessionAwareMessageL jmsReply.setJMSCorrelationID(jmsMessage.getJMSMessageID()); } MessageProducer producer = session.createProducer(destination); - producer.setTimeToLive(this.replyTimeToLive); - producer.setPriority(this.replyPriority); - producer.setDeliveryMode(this.replyDeliveryMode); try { - producer.send(jmsReply); + if (this.explicitQosEnabledForReplies) { + producer.send(jmsReply, + this.replyDeliveryMode, this.replyPriority, this.replyTimeToLive); + } + else { + producer.send(jmsReply); + } } finally { producer.close(); diff --git a/org.springframework.integration.jms/src/main/java/org/springframework/integration/jms/config/JmsMessageDrivenEndpointParser.java b/org.springframework.integration.jms/src/main/java/org/springframework/integration/jms/config/JmsMessageDrivenEndpointParser.java index f73d59358e..d74c310099 100644 --- a/org.springframework.integration.jms/src/main/java/org/springframework/integration/jms/config/JmsMessageDrivenEndpointParser.java +++ b/org.springframework.integration.jms/src/main/java/org/springframework/integration/jms/config/JmsMessageDrivenEndpointParser.java @@ -47,6 +47,8 @@ public class JmsMessageDrivenEndpointParser extends AbstractSingleBeanDefinition private static final String REPLY_DELIVERY_PERSISTENT = "reply-delivery-persistent"; + private static final String EXPLICIT_QOS_ENABLED_FOR_REPLIES = "explicit-qos-enabled-for-replies"; + private static String[] containerAttributes = new String[] { JmsAdapterParserUtils.CONNECTION_FACTORY_PROPERTY, @@ -176,6 +178,7 @@ public class JmsMessageDrivenEndpointParser extends AbstractSingleBeanDefinition IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, REPLY_TIME_TO_LIVE); IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, REPLY_PRIORITY); IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, REPLY_DELIVERY_PERSISTENT); + IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, EXPLICIT_QOS_ENABLED_FOR_REPLIES); } else { IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "channel", "requestChannel"); diff --git a/org.springframework.integration.jms/src/main/resources/org/springframework/integration/jms/config/spring-integration-jms-2.0.xsd b/org.springframework.integration.jms/src/main/resources/org/springframework/integration/jms/config/spring-integration-jms-2.0.xsd index 1406ff5eae..e38d2621e7 100644 --- a/org.springframework.integration.jms/src/main/resources/org/springframework/integration/jms/config/spring-integration-jms-2.0.xsd +++ b/org.springframework.integration.jms/src/main/resources/org/springframework/integration/jms/config/spring-integration-jms-2.0.xsd @@ -522,6 +522,7 @@ + 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 b7ba2240a7..de26d90c29 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 @@ -274,6 +274,17 @@ public class JmsInboundGatewayParserTests { assertEquals(12345L, accessor.getPropertyValue("replyTimeToLive")); assertEquals(7, accessor.getPropertyValue("replyPriority")); assertEquals(DeliveryMode.NON_PERSISTENT, accessor.getPropertyValue("replyDeliveryMode")); + assertEquals(true, accessor.getPropertyValue("explicitQosEnabledForReplies")); + } + + @Test + public void replyQosPropertiesDisabledByDefault() { + ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext( + "inboundGatewayDefault.xml", this.getClass()); + JmsMessageDrivenEndpoint gateway = context.getBean("gateway", JmsMessageDrivenEndpoint.class); + DirectFieldAccessor accessor = new DirectFieldAccessor( + new DirectFieldAccessor(gateway).getPropertyValue("listener")); + assertEquals(false, accessor.getPropertyValue("explicitQosEnabledForReplies")); } } diff --git a/org.springframework.integration.jms/src/test/java/org/springframework/integration/jms/config/inboundGatewayDefault.xml b/org.springframework.integration.jms/src/test/java/org/springframework/integration/jms/config/inboundGatewayDefault.xml new file mode 100644 index 0000000000..1dc63b9d7c --- /dev/null +++ b/org.springframework.integration.jms/src/test/java/org/springframework/integration/jms/config/inboundGatewayDefault.xml @@ -0,0 +1,29 @@ + + + + + + + + + + + + + + + + + + diff --git a/org.springframework.integration.jms/src/test/java/org/springframework/integration/jms/config/inboundGatewayWithReplyQos.xml b/org.springframework.integration.jms/src/test/java/org/springframework/integration/jms/config/inboundGatewayWithReplyQos.xml index 915a728b2b..18ca2dd35f 100644 --- a/org.springframework.integration.jms/src/test/java/org/springframework/integration/jms/config/inboundGatewayWithReplyQos.xml +++ b/org.springframework.integration.jms/src/test/java/org/springframework/integration/jms/config/inboundGatewayWithReplyQos.xml @@ -19,7 +19,8 @@ request-channel="requestChannel" reply-time-to-live="12345" reply-priority="7" - reply-delivery-persistent="false"/> + reply-delivery-persistent="false" + explicit-qos-enabled-for-replies="true"/>