INT-580
This commit is contained in:
@@ -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.
|
||||
* <p>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.
|
||||
* <p>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.
|
||||
* <p>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.
|
||||
* <p>This implementation first checks the JMS Reply-To {@link Destination}
|
||||
* of the supplied request message; if that is not <code>null</code> it is
|
||||
* returned; if it is <code>null</code>, then the configured
|
||||
* {@link #resolveDefaultReplyDestination default reply destination}
|
||||
* is returned; if this too is <code>null</code>, 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 <code>null</code>)
|
||||
* @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;
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user