diff --git a/support/src/main/java/org/springframework/ws/transport/jms/JmsSenderConnection.java b/support/src/main/java/org/springframework/ws/transport/jms/JmsSenderConnection.java index d6528fe5..7d565616 100644 --- a/support/src/main/java/org/springframework/ws/transport/jms/JmsSenderConnection.java +++ b/support/src/main/java/org/springframework/ws/transport/jms/JmsSenderConnection.java @@ -32,6 +32,7 @@ import javax.jms.MessageConsumer; import javax.jms.MessageProducer; import javax.jms.Session; import javax.jms.TemporaryQueue; +import javax.jms.TemporaryTopic; import javax.jms.TextMessage; import org.springframework.jms.connection.ConnectionFactoryUtils; @@ -214,7 +215,14 @@ public class JmsSenderConnection extends AbstractSenderConnection { protected void onReceiveBeforeRead() throws IOException { MessageConsumer messageConsumer = null; try { - messageConsumer = session.createConsumer(responseDestination); + if (responseDestination instanceof TemporaryQueue || responseDestination instanceof TemporaryTopic) { + messageConsumer = session.createConsumer(responseDestination); + } + else { + String messageId = requestMessage.getJMSMessageID().replaceAll("'", "''"); + String messageSelector = "JMSCorrelationID = '" + messageId + "'"; + messageConsumer = session.createConsumer(responseDestination, messageSelector); + } Message message = receiveTimeout >= 0 ? messageConsumer.receive(receiveTimeout) : messageConsumer.receive(); if (message instanceof BytesMessage || message instanceof TextMessage) { responseMessage = message; @@ -237,6 +245,14 @@ public class JmsSenderConnection extends AbstractSenderConnection { // ignore } } + else if (responseDestination instanceof TemporaryTopic) { + try { + ((TemporaryTopic) responseDestination).delete(); + } + catch (JMSException e) { + // ignore + } + } } } diff --git a/support/src/test/java/org/springframework/ws/transport/jms/JmsIntegrationTest.java b/support/src/test/java/org/springframework/ws/transport/jms/JmsIntegrationTest.java index f5a0e7a0..1561e404 100644 --- a/support/src/test/java/org/springframework/ws/transport/jms/JmsIntegrationTest.java +++ b/support/src/test/java/org/springframework/ws/transport/jms/JmsIntegrationTest.java @@ -16,13 +16,13 @@ package org.springframework.ws.transport.jms; +import org.custommonkey.xmlunit.XMLAssert; + import org.springframework.test.AbstractDependencyInjectionSpringContextTests; import org.springframework.ws.client.core.WebServiceTemplate; import org.springframework.xml.transform.StringResult; import org.springframework.xml.transform.StringSource; -import org.custommonkey.xmlunit.XMLAssert; - public class JmsIntegrationTest extends AbstractDependencyInjectionSpringContextTests { private WebServiceTemplate webServiceTemplate; @@ -35,11 +35,23 @@ public class JmsIntegrationTest extends AbstractDependencyInjectionSpringContext this.webServiceTemplate = webServiceTemplate; } - public void testJmsTransport() throws Exception { + protected void onTearDown() throws Exception { + applicationContext.close(); + setDirty(); + } + + public void testTemporaryQueue() throws Exception { String content = ""; StringResult result = new StringResult(); webServiceTemplate.sendSourceAndReceiveToResult(new StringSource(content), result); XMLAssert.assertXMLEqual("Invalid content received", content, result.toString()); - applicationContext.close(); + } + + public void testPermanentQueue() throws Exception { + String url = "jms:RequestQueue?deliveryMode=NON_PERSISTENT;replyToName=ResponseQueue"; + String content = ""; + StringResult result = new StringResult(); + webServiceTemplate.sendSourceAndReceiveToResult(url, new StringSource(content), result); + XMLAssert.assertXMLEqual("Invalid content received", content, result.toString()); } }