diff --git a/spring-integration-jms/pom.xml b/spring-integration-jms/pom.xml index 6b234961cc..7f23ba7858 100644 --- a/spring-integration-jms/pom.xml +++ b/spring-integration-jms/pom.xml @@ -35,6 +35,11 @@ junit junit + + org.mockito + mockito-all + test + org.springframework spring-test diff --git a/spring-integration-jms/src/main/java/org/springframework/integration/jms/AbstractJmsTemplateBasedAdapter.java b/spring-integration-jms/src/main/java/org/springframework/integration/jms/AbstractJmsTemplateBasedAdapter.java index 4c33f0b18a..57f0ec67ec 100644 --- a/spring-integration-jms/src/main/java/org/springframework/integration/jms/AbstractJmsTemplateBasedAdapter.java +++ b/spring-integration-jms/src/main/java/org/springframework/integration/jms/AbstractJmsTemplateBasedAdapter.java @@ -20,7 +20,7 @@ import javax.jms.ConnectionFactory; import javax.jms.DeliveryMode; import javax.jms.Destination; -import org.springframework.beans.factory.InitializingBean; +import org.springframework.integration.context.IntegrationObjectSupport; import org.springframework.jms.core.JmsTemplate; import org.springframework.jms.support.converter.MessageConverter; import org.springframework.jms.support.converter.SimpleMessageConverter; @@ -33,7 +33,7 @@ import org.springframework.util.Assert; * @author Mark Fisher * @author Oleg Zhurakousky */ -public abstract class AbstractJmsTemplateBasedAdapter implements InitializingBean { +public abstract class AbstractJmsTemplateBasedAdapter extends IntegrationObjectSupport { private volatile boolean extractPayload = true; @@ -181,7 +181,7 @@ public abstract class AbstractJmsTemplateBasedAdapter implements InitializingBea return this.jmsTemplate; } - public void afterPropertiesSet() { + public void onInit() { synchronized (this.initializationMonitor) { if (this.initialized) { return; diff --git a/spring-integration-jms/src/main/java/org/springframework/integration/jms/ChannelPublishingJmsMessageListener.java b/spring-integration-jms/src/main/java/org/springframework/integration/jms/ChannelPublishingJmsMessageListener.java index c1087753a1..f3d55274b5 100644 --- a/spring-integration-jms/src/main/java/org/springframework/integration/jms/ChannelPublishingJmsMessageListener.java +++ b/spring-integration-jms/src/main/java/org/springframework/integration/jms/ChannelPublishingJmsMessageListener.java @@ -16,6 +16,8 @@ package org.springframework.integration.jms; +import java.util.Map; + import javax.jms.DeliveryMode; import javax.jms.Destination; import javax.jms.InvalidDestinationException; @@ -67,6 +69,10 @@ public class ChannelPublishingJmsMessageListener extends AbstractMessagingGatewa private volatile DestinationResolver destinationResolver = new DynamicDestinationResolver(); private volatile JmsHeaderMapper headerMapper = new DefaultJmsHeaderMapper(); + + public String getComponentType(){ + return "jms:inbound-gateway"; + } /** * Specify whether a JMS reply Message is expected. @@ -208,10 +214,15 @@ public class ChannelPublishingJmsMessageListener extends AbstractMessagingGatewa super.onInit(); } + @SuppressWarnings("unchecked") public void onMessage(javax.jms.Message jmsMessage, Session session) throws JMSException { Object object = this.messageConverter.fromMessage(jmsMessage); + + Map headers = (Map) headerMapper.toHeaders(jmsMessage); Message requestMessage = (object instanceof Message) ? - (Message) object : MessageBuilder.withPayload(object).build(); + MessageBuilder.fromMessage((Message) object).copyHeaders(headers).build() : + MessageBuilder.withPayload(object).copyHeaders(headers).build(); + this.writeMessageHistory(requestMessage, this); if (!this.expectReply) { this.send(requestMessage); } diff --git a/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsDestinationPollingSource.java b/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsDestinationPollingSource.java index efa09f61cf..f7de52ce46 100644 --- a/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsDestinationPollingSource.java +++ b/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsDestinationPollingSource.java @@ -53,13 +53,16 @@ public class JmsDestinationPollingSource extends AbstractJmsTemplateBasedAdapter super(connectionFactory, destinationName); } - + public String getComponentType(){ + return "jms:inbound-channel-adapter"; + } /** * Specify a JMS Message Selector expression to use when receiving Messages. */ public void setMessageSelector(String messageSelector) { this.messageSelector = messageSelector; } + /** * Will receive JMS {@link javax.jms.Message} converting and returning it as * Spring Integration(SI) {@link Message}. @@ -86,6 +89,7 @@ public class JmsDestinationPollingSource extends AbstractJmsTemplateBasedAdapter } else { convertedMessage = MessageBuilder.withPayload(convertedObject).build(); } + this.writeMessageHistory(convertedMessage, this); } catch (Exception e) { throw new MessagingException(e.getMessage(), e); } diff --git a/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsMessageDrivenEndpoint.java b/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsMessageDrivenEndpoint.java index bacb44dbc3..212aca2cb4 100644 --- a/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsMessageDrivenEndpoint.java +++ b/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsMessageDrivenEndpoint.java @@ -53,6 +53,7 @@ public class JmsMessageDrivenEndpoint extends AbstractEndpoint implements Dispos if (!this.listenerContainer.isActive()) { this.listenerContainer.afterPropertiesSet(); } + listener.setComponentName(this.getComponentName()); } @Override diff --git a/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsOutboundGateway.java b/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsOutboundGateway.java index c309ee2635..39becd0eaa 100644 --- a/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsOutboundGateway.java +++ b/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsOutboundGateway.java @@ -88,6 +88,10 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler { private volatile boolean initialized; private final Object initializationMonitor = new Object(); + + public String getComponentType(){ + return "jms:outbound-gateway"; + } /** diff --git a/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsSendingMessageHandler.java b/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsSendingMessageHandler.java index c9169d23bd..20d5cc4aed 100644 --- a/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsSendingMessageHandler.java +++ b/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsSendingMessageHandler.java @@ -35,7 +35,9 @@ public class JmsSendingMessageHandler extends AbstractJmsTemplateBasedAdapter im private volatile int order = Ordered.LOWEST_PRECEDENCE; - + public String getComponentType(){ + return "jms:outbound-channel-adapter"; + } public JmsSendingMessageHandler(JmsTemplate jmsTemplate) { super(jmsTemplate); } @@ -60,6 +62,7 @@ public class JmsSendingMessageHandler extends AbstractJmsTemplateBasedAdapter im if (message == null) { throw new IllegalArgumentException("message must not be null"); } + this.writeMessageHistory(message, this); this.getJmsTemplate().convertAndSend(message, new MessagePostProcessor() { public javax.jms.Message postProcessMessage(javax.jms.Message jmsMessage) throws JMSException { diff --git a/spring-integration-jms/src/main/java/org/springframework/integration/jms/config/JmsInboundChannelAdapterParser.java b/spring-integration-jms/src/main/java/org/springframework/integration/jms/config/JmsInboundChannelAdapterParser.java index a3e2c59d4e..b952fbfa06 100644 --- a/spring-integration-jms/src/main/java/org/springframework/integration/jms/config/JmsInboundChannelAdapterParser.java +++ b/spring-integration-jms/src/main/java/org/springframework/integration/jms/config/JmsInboundChannelAdapterParser.java @@ -48,6 +48,10 @@ public class JmsInboundChannelAdapterParser extends AbstractPollingInboundChanne Object source = parserContext.extractSource(element); BeanDefinitionBuilder builder = BeanDefinitionBuilder.genericBeanDefinition( "org.springframework.integration.jms.JmsDestinationPollingSource"); + String componentName = this.resolveId(element, builder.getBeanDefinition(), parserContext); + if (StringUtils.hasText(componentName)){ + builder.addPropertyValue("componentName", componentName); + } String jmsTemplate = element.getAttribute(JmsAdapterParserUtils.JMS_TEMPLATE_ATTRIBUTE); String destination = element.getAttribute(JmsAdapterParserUtils.DESTINATION_ATTRIBUTE); String destinationName = element.getAttribute(JmsAdapterParserUtils.DESTINATION_NAME_ATTRIBUTE); diff --git a/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/JmsMessageHistoryTests.java b/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/JmsMessageHistoryTests.java new file mode 100644 index 0000000000..22f1957a5b --- /dev/null +++ b/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/JmsMessageHistoryTests.java @@ -0,0 +1,188 @@ +/* + * Copyright 2002-2010 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.integration.jms.config; +import static junit.framework.Assert.assertEquals; + +import java.util.Iterator; +import java.util.Map; +import java.util.StringTokenizer; + +import org.junit.Test; +import org.mockito.Mockito; +import org.springframework.context.ConfigurableApplicationContext; +import org.springframework.context.support.ClassPathXmlApplicationContext; +import org.springframework.integration.channel.DirectChannel; +import org.springframework.integration.channel.PollableChannel; +import org.springframework.integration.channel.SubscribableChannel; +import org.springframework.integration.context.NamedComponent; +import org.springframework.integration.core.Message; +import org.springframework.integration.core.MessageHeaders; +import org.springframework.integration.core.MessagingException; +import org.springframework.integration.history.MessageHistory; +import org.springframework.integration.history.MessageHistoryEvent; +import org.springframework.integration.jms.DefaultJmsHeaderMapper; +import org.springframework.integration.message.MessageDeliveryException; +import org.springframework.integration.message.MessageHandler; +import org.springframework.integration.message.MessageHandlingException; +import org.springframework.integration.message.MessageRejectedException; +import org.springframework.integration.message.StringMessage; +import org.springframework.util.StringUtils; + +/** + * @author Oleg Zhurakousky + * + */ +public class JmsMessageHistoryTests { + + @SuppressWarnings("unchecked") + @Test + public void testInboundAdapter() throws Exception{ + ActiveMqTestUtils.prepare(); + ConfigurableApplicationContext applicationContext = new ClassPathXmlApplicationContext("MessageHistoryTests-context.xml", JmsMessageHistoryTests.class); + SampleGateway gateway = applicationContext.getBean("sampleGateway", SampleGateway.class); + PollableChannel jmsInputChannel = applicationContext.getBean("jmsInputChannel", PollableChannel.class); + gateway.send("hello"); + Message message = (Message) jmsInputChannel.receive(5000); + Iterator historyIterator = message.getHeaders().getHistory().iterator(); + MessageHistoryEvent event = historyIterator.next(); + assertEquals("jms:inbound-channel-adapter", event.getType()); + assertEquals("sampleJmsInboundAdapter", event.getName()); + event = historyIterator.next(); + assertEquals("queue-channel", event.getType()); + assertEquals("jmsInputChannel", event.getName()); + } + @SuppressWarnings("unchecked") + @Test + public void testWithHeaderMapperPropagatingOutboundHistory() throws Exception{ + ActiveMqTestUtils.prepare(); + ConfigurableApplicationContext applicationContext = new ClassPathXmlApplicationContext("MessageHistoryTests-withHeaderMapper.xml", JmsMessageHistoryTests.class); + DirectChannel input = applicationContext.getBean("outbound-channel", DirectChannel.class); + PollableChannel jmsInputChannel = applicationContext.getBean("jmsInputChannel", PollableChannel.class); + input.send(new StringMessage("hello")); + Message message = (Message) jmsInputChannel.receive(50000); + Iterator historyIterator = message.getHeaders().getHistory().iterator(); + MessageHistoryEvent event = historyIterator.next(); + assertEquals("channel", event.getType()); + assertEquals("outbound-channel", event.getName()); + event = historyIterator.next(); + assertEquals("jms:outbound-channel-adapter", event.getType()); + assertEquals("jmsOutbound", event.getName()); + event = historyIterator.next(); + assertEquals("jms:inbound-channel-adapter", event.getType()); + assertEquals("sampleJmsInboundAdapter", event.getName()); + event = historyIterator.next(); + assertEquals("queue-channel", event.getType()); + assertEquals("jmsInputChannel", event.getName()); + } + @SuppressWarnings("unchecked") + @Test + public void testWithHeaderMapperPropagatingOutboundHistoryWithGateways() throws Exception{ + ActiveMqTestUtils.prepare(); + ConfigurableApplicationContext applicationContext = new ClassPathXmlApplicationContext("MessageHistoryTests-gateways.xml", JmsMessageHistoryTests.class); + SampleGateway gateway = applicationContext.getBean("sampleGateway", SampleGateway.class); + SubscribableChannel inboundJmsChannel = applicationContext.getBean("inbound-jms-channel", SubscribableChannel.class); + MessageHandler handler = new MessageHandler() { + public void handleMessage(Message message) + throws MessageRejectedException, MessageHandlingException, + MessageDeliveryException { + Iterator historyIterator = message.getHeaders().getHistory().iterator(); + MessageHistoryEvent event = historyIterator.next(); + assertEquals("gateway", event.getType()); + assertEquals("sampleGateway", event.getName()); + event = historyIterator.next(); + assertEquals("pub-sub-channel", event.getType()); + assertEquals("channel-a", event.getName()); + event = historyIterator.next(); + assertEquals("jms:outbound-gateway", event.getType()); + assertEquals("jmsOutbound", event.getName()); + event = historyIterator.next(); + assertEquals("jms:inbound-gateway", event.getType()); + assertEquals("jmsInbound", event.getName()); + event = historyIterator.next(); + assertEquals("pub-sub-channel", event.getType()); + assertEquals("inbound-jms-channel", event.getName()); + event = historyIterator.next(); + assertEquals("service-activator", event.getType()); + assertEquals("sampleService-a", event.getName()); + } + }; + handler = Mockito.spy(handler); + inboundJmsChannel.subscribe(handler); + gateway.echo("hello"); + Mockito.verify(handler, Mockito.times(1)).handleMessage(Mockito.any(Message.class)); + } + + public static interface SampleGateway{ + public void send(String value); + public Message echo(String value); + } + + public static class SampleService{ + public Message echoMessage(String value){ + return new StringMessage(value); + } + } + + public static class SampleHeaderMapper extends DefaultJmsHeaderMapper { + + + public void fromHeaders(MessageHeaders headers, javax.jms.Message jmsMessage){ + super.fromHeaders(headers, jmsMessage); + String messageHistory = headers.getHistory().toString(); + try { + jmsMessage.setStringProperty("outbound_history", messageHistory); + } catch (Exception e) { + throw new MessagingException("Problem setting JMS properties", e); + } + } + + public Map toHeaders(javax.jms.Message jmsMessage){ + Map headers = super.toHeaders(jmsMessage); + MessageHistory history = new MessageHistory(); + String outboundHistory = (String) headers.get("outbound_history"); + StringTokenizer tok = new StringTokenizer(outboundHistory, ",[] "); + while (tok.hasMoreTokens()) { + String historyItem = tok.nextToken(); + String[] parsedHistory = StringUtils.split(historyItem, "@"); + String type = null; + String name = historyItem; + if (parsedHistory != null){ + name = parsedHistory[1]; + type = parsedHistory[0]; + } + history.addEvent(new SampleComponent(name, type)); + } + headers.put(MessageHeaders.HISTORY, history); + headers.remove("outbound_history"); + return headers; + } + } + + public static class SampleComponent implements NamedComponent{ + private String name; + private String type; + public SampleComponent(String name, String type){ + this.name = name; + this.type = type; + } + public String getComponentName() { + return name; + } + public String getComponentType() { + return type; + } + } +} diff --git a/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/MessageHistoryTests-context.xml b/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/MessageHistoryTests-context.xml new file mode 100644 index 0000000000..f06c224314 --- /dev/null +++ b/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/MessageHistoryTests-context.xml @@ -0,0 +1,46 @@ + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + diff --git a/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/MessageHistoryTests-gateways.xml b/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/MessageHistoryTests-gateways.xml new file mode 100644 index 0000000000..4256cf881e --- /dev/null +++ b/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/MessageHistoryTests-gateways.xml @@ -0,0 +1,44 @@ + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + diff --git a/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/MessageHistoryTests-withHeaderMapper.xml b/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/MessageHistoryTests-withHeaderMapper.xml new file mode 100644 index 0000000000..73bb07ca97 --- /dev/null +++ b/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/MessageHistoryTests-withHeaderMapper.xml @@ -0,0 +1,42 @@ + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + +