From 004cee34a4f567150b2c1ffd496f532aac78bbfb Mon Sep 17 00:00:00 2001 From: Mark Fisher Date: Thu, 19 Feb 2009 18:38:36 +0000 Subject: [PATCH] INT-580 --- .../ChannelPublishingJmsMessageListener.java | 127 ++++++++++++++++-- 1 file changed, 115 insertions(+), 12 deletions(-) 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 061180376a..06a516c855 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 @@ -17,6 +17,7 @@ package org.springframework.integration.jms; import javax.jms.Destination; +import javax.jms.InvalidDestinationException; import javax.jms.JMSException; import javax.jms.MessageProducer; import javax.jms.Session; @@ -25,11 +26,13 @@ import org.springframework.beans.factory.InitializingBean; import org.springframework.integration.channel.MessageChannelTemplate; import org.springframework.integration.core.Message; import org.springframework.integration.core.MessageChannel; -import org.springframework.integration.core.MessagingException; import org.springframework.integration.message.MessageBuilder; import org.springframework.integration.message.MessageDeliveryException; import org.springframework.jms.listener.SessionAwareMessageListener; import org.springframework.jms.support.converter.MessageConverter; +import org.springframework.jms.support.destination.DestinationResolver; +import org.springframework.jms.support.destination.DynamicDestinationResolver; +import org.springframework.util.Assert; /** * JMS MessageListener that converts a JMS Message into a Spring Integration @@ -38,6 +41,7 @@ import org.springframework.jms.support.converter.MessageConverter; * and convert that into a JMS reply. * * @author Mark Fisher + * @author Juergen Hoeller */ public class ChannelPublishingJmsMessageListener implements SessionAwareMessageListener, InitializingBean { @@ -49,7 +53,9 @@ public class ChannelPublishingJmsMessageListener implements SessionAwareMessageL private volatile boolean extractReplyPayload = true; - private volatile Destination defaultReplyDestination; + private volatile Object defaultReplyDestination; + + private volatile DestinationResolver destinationResolver = new DynamicDestinationResolver(); private volatile JmsHeaderMapper headerMapper; @@ -89,13 +95,51 @@ public class ChannelPublishingJmsMessageListener implements SessionAwareMessageL } /** - * Specify the default reply Destination. If a request Message does not provide - * a 'JMSReplyTo' property, replies will be sent to this by default. + * Set the default reply destination to send reply messages to. This will + * be applied in case of a request message that does not carry a + * "JMSReplyTo" field. */ public void setDefaultReplyDestination(Destination defaultReplyDestination) { this.defaultReplyDestination = defaultReplyDestination; } + /** + * Set the name of the default reply queue to send reply messages to. + * This will be applied in case of a request message that does not carry a + * "JMSReplyTo" field. + *

Alternatively, specify a JMS Destination object as "defaultReplyDestination". + * @see #setDestinationResolver + * @see #setDefaultReplyDestination(javax.jms.Destination) + */ + public void setDefaultReplyQueueName(String destinationName) { + this.defaultReplyDestination = new DestinationNameHolder(destinationName, false); + } + + /** + * Set the name of the default reply topic to send reply messages to. + * This will be applied in case of a request message that does not carry a + * "JMSReplyTo" field. + *

Alternatively, specify a JMS Destination object as "defaultReplyDestination". + * @see #setDestinationResolver + * @see #setDefaultReplyDestination(javax.jms.Destination) + */ + public void setDefaultReplyTopicName(String destinationName) { + this.defaultReplyDestination = new DestinationNameHolder(destinationName, true); + } + + /** + * Set the DestinationResolver that should be used to resolve reply + * destination names for this listener. + *

The default resolver is a DynamicDestinationResolver. Specify a + * JndiDestinationResolver for resolving destination names as JNDI locations. + * @see org.springframework.jms.support.destination.DynamicDestinationResolver + * @see org.springframework.jms.support.destination.JndiDestinationResolver + */ + public void setDestinationResolver(DestinationResolver destinationResolver) { + Assert.notNull(destinationResolver, "destinationResolver must not be null"); + this.destinationResolver = destinationResolver; + } + /** * Provide a {@link MessageConverter} implementation to use when * converting between JMS Messages and Spring Integration Messages. @@ -164,14 +208,7 @@ public class ChannelPublishingJmsMessageListener implements SessionAwareMessageL else { Message replyMessage = this.channelTemplate.sendAndReceive(requestMessage); if (replyMessage != null) { - Destination destination = jmsMessage.getJMSReplyTo(); - if (destination == null) { - destination = this.defaultReplyDestination; - } - if (destination == null) { - throw new MessagingException(replyMessage, "Unable to send JMS reply. The request Message " - + "has no 'JMSReplyTo' property, and this listener has no 'defaultReplyDestination'."); - } + Destination destination = this.getReplyDestination(jmsMessage, session); javax.jms.Message jmsReply = this.messageConverter.toMessage(replyMessage, session); if (jmsReply.getJMSCorrelationID() == null) { jmsReply.setJMSCorrelationID(jmsMessage.getJMSMessageID()); @@ -182,4 +219,70 @@ public class ChannelPublishingJmsMessageListener implements SessionAwareMessageL } } + /** + * Determine a reply destination for the given message. + *

This implementation first checks the JMS Reply-To {@link Destination} + * of the supplied request message; if that is not null it is + * returned; if it is null, then the configured + * {@link #resolveDefaultReplyDestination default reply destination} + * is returned; if this too is null, then an + * {@link InvalidDestinationException} is thrown. + * @param request the original incoming JMS message + * @param session the JMS Session to operate on + * @return the reply destination (never null) + * @throws JMSException if thrown by JMS API methods + * @throws InvalidDestinationException if no {@link Destination} can be determined + * @see #setDefaultReplyDestination + * @see javax.jms.Message#getJMSReplyTo() + */ + private Destination getReplyDestination(javax.jms.Message request, Session session) throws JMSException { + Destination replyTo = request.getJMSReplyTo(); + if (replyTo == null) { + replyTo = resolveDefaultReplyDestination(session); + if (replyTo == null) { + throw new InvalidDestinationException("Cannot determine reply destination: " + + "Request message does not contain reply-to destination, and no default reply destination set."); + } + } + return replyTo; + } + + /** + * Resolve the default reply destination into a JMS {@link Destination}, using this + * listener's {@link DestinationResolver} in case of a destination name. + * @return the located {@link Destination} + * @throws javax.jms.JMSException if resolution failed + * @see #setDefaultReplyDestination + * @see #setDefaultReplyQueueName + * @see #setDefaultReplyTopicName + * @see #setDestinationResolver + */ + private Destination resolveDefaultReplyDestination(Session session) throws JMSException { + if (this.defaultReplyDestination instanceof Destination) { + return (Destination) this.defaultReplyDestination; + } + if (this.defaultReplyDestination instanceof DestinationNameHolder) { + DestinationNameHolder nameHolder = (DestinationNameHolder) this.defaultReplyDestination; + return this.destinationResolver.resolveDestinationName(session, nameHolder.name, nameHolder.isTopic); + } + return null; + } + + + /** + * Internal class combining a destination name + * and its target destination type (queue or topic). + */ + private static class DestinationNameHolder { + + public final String name; + + public final boolean isTopic; + + public DestinationNameHolder(String name, boolean isTopic) { + this.name = name; + this.isTopic = isTopic; + } + } + }