From f45437801301d2b1aa285d66732f77aca75c30a5 Mon Sep 17 00:00:00 2001 From: Mark Fisher Date: Mon, 13 Sep 2010 13:10:53 -0400 Subject: [PATCH 1/7] INT-1412 added async send methods to AsyncMessagingOperations and AsyncMessagingTemplate --- .../core/AsyncMessagingOperations.java | 12 +++ .../core/AsyncMessagingTemplate.java | 48 ++++++++++ .../core/AsyncMessagingTemplateTests.java | 92 +++++++++++++++++++ 3 files changed, 152 insertions(+) diff --git a/spring-integration-core/src/main/java/org/springframework/integration/core/AsyncMessagingOperations.java b/spring-integration-core/src/main/java/org/springframework/integration/core/AsyncMessagingOperations.java index 7b29161f94..4d735344be 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/core/AsyncMessagingOperations.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/core/AsyncMessagingOperations.java @@ -27,6 +27,18 @@ import org.springframework.integration.MessageChannel; */ public interface AsyncMessagingOperations { + Future asyncSend(Message message); + + Future asyncSend(MessageChannel channel, Message message); + + Future asyncSend(String channelName, Message message); + + Future asyncConvertAndSend(Object message); + + Future asyncConvertAndSend(MessageChannel channel, Object message); + + Future asyncConvertAndSend(String channelName, Object message); + Future> asyncReceive(); Future> asyncReceive(PollableChannel channel); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/core/AsyncMessagingTemplate.java b/spring-integration-core/src/main/java/org/springframework/integration/core/AsyncMessagingTemplate.java index 12fc4fb0cc..6118cea71d 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/core/AsyncMessagingTemplate.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/core/AsyncMessagingTemplate.java @@ -42,6 +42,54 @@ public class AsyncMessagingTemplate extends MessagingTemplate implements AsyncMe (AsyncTaskExecutor) executor : new TaskExecutorAdapter(executor); } + public Future asyncSend(final Message message) { + return this.executor.submit(new Runnable() { + public void run() { + send(message); + } + }); + } + + public Future asyncSend(final MessageChannel channel, final Message message) { + return this.executor.submit(new Runnable() { + public void run() { + send(channel, message); + } + }); + } + + public Future asyncSend(final String channelName, final Message message) { + return this.executor.submit(new Runnable() { + public void run() { + send(channelName, message); + } + }); + } + + public Future asyncConvertAndSend(final Object object) { + return this.executor.submit(new Runnable() { + public void run() { + convertAndSend(object); + } + }); + } + + public Future asyncConvertAndSend(final MessageChannel channel, final Object object) { + return this.executor.submit(new Runnable() { + public void run() { + convertAndSend(channel, object); + } + }); + } + + public Future asyncConvertAndSend(final String channelName, final Object object) { + return this.executor.submit(new Runnable() { + public void run() { + convertAndSend(channelName, object); + } + }); + } + public Future> asyncReceive() { return this.executor.submit(new Callable>() { public Message call() throws Exception { diff --git a/spring-integration-core/src/test/java/org/springframework/integration/core/AsyncMessagingTemplateTests.java b/spring-integration-core/src/test/java/org/springframework/integration/core/AsyncMessagingTemplateTests.java index 5458cbcbed..cb31007ac5 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/core/AsyncMessagingTemplateTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/core/AsyncMessagingTemplateTests.java @@ -18,6 +18,7 @@ package org.springframework.integration.core; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertNull; import static org.junit.Assert.assertTrue; import static org.junit.Assert.fail; @@ -47,6 +48,97 @@ import org.springframework.util.Assert; */ public class AsyncMessagingTemplateTests { + @Test + public void asyncSendWithDefaultChannel() throws Exception { + QueueChannel channel = new QueueChannel(); + AsyncMessagingTemplate template = new AsyncMessagingTemplate(); + template.setDefaultChannel(channel); + Message message = MessageBuilder.withPayload("test").build(); + Future future = template.asyncSend(message); + assertNull(future.get(1000, TimeUnit.MILLISECONDS)); + Message result = channel.receive(0); + assertEquals(message, result); + } + + @Test + public void asyncSendWithExplicitChannel() throws Exception { + QueueChannel channel = new QueueChannel(); + AsyncMessagingTemplate template = new AsyncMessagingTemplate(); + Message message = MessageBuilder.withPayload("test").build(); + Future future = template.asyncSend(channel, message); + assertNull(future.get(1000, TimeUnit.MILLISECONDS)); + Message result = channel.receive(0); + assertEquals(message, result); + } + + @Test + public void asyncSendWithResolvedChannel() throws Exception { + StaticApplicationContext context = new StaticApplicationContext(); + context.registerSingleton("testChannel", QueueChannel.class); + context.refresh(); + QueueChannel channel = context.getBean("testChannel", QueueChannel.class); + AsyncMessagingTemplate template = new AsyncMessagingTemplate(); + template.setBeanFactory(context); + Message message = MessageBuilder.withPayload("test").build(); + Future future = template.asyncSend("testChannel", message); + assertNull(future.get(1000, TimeUnit.MILLISECONDS)); + Message result = channel.receive(0); + assertEquals(message, result); + } + + @Test(expected = TimeoutException.class) + public void asyncSendWithTimeoutException() throws Exception { + QueueChannel channel = new QueueChannel(1); + channel.send(MessageBuilder.withPayload("blocker").build()); + AsyncMessagingTemplate template = new AsyncMessagingTemplate(); + Future result = template.asyncSend(channel, MessageBuilder.withPayload("test").build()); + result.get(100, TimeUnit.MILLISECONDS); + } + + @Test + public void asyncConvertAndSendWithDefaultChannel() throws Exception { + QueueChannel channel = new QueueChannel(); + AsyncMessagingTemplate template = new AsyncMessagingTemplate(); + template.setDefaultChannel(channel); + Future future = template.asyncConvertAndSend("test"); + assertNull(future.get(1000, TimeUnit.MILLISECONDS)); + Message result = channel.receive(0); + assertEquals("test", result.getPayload()); + } + + @Test + public void asyncConvertAndSendWithExplicitChannel() throws Exception { + QueueChannel channel = new QueueChannel(); + AsyncMessagingTemplate template = new AsyncMessagingTemplate(); + Future future = template.asyncConvertAndSend(channel, "test"); + assertNull(future.get(1000, TimeUnit.MILLISECONDS)); + Message result = channel.receive(0); + assertEquals("test", result.getPayload()); + } + + @Test + public void asyncConvertAndSendWithResolvedChannel() throws Exception { + StaticApplicationContext context = new StaticApplicationContext(); + context.registerSingleton("testChannel", QueueChannel.class); + context.refresh(); + QueueChannel channel = context.getBean("testChannel", QueueChannel.class); + AsyncMessagingTemplate template = new AsyncMessagingTemplate(); + template.setBeanFactory(context); + Future future = template.asyncConvertAndSend("testChannel", "test"); + assertNull(future.get(1000, TimeUnit.MILLISECONDS)); + Message result = channel.receive(0); + assertEquals("test", result.getPayload()); + } + + @Test(expected = TimeoutException.class) + public void asyncConvertAndSendWithTimeoutException() throws Exception { + QueueChannel channel = new QueueChannel(1); + channel.send(MessageBuilder.withPayload("blocker").build()); + AsyncMessagingTemplate template = new AsyncMessagingTemplate(); + Future result = template.asyncConvertAndSend(channel, "test"); + result.get(100, TimeUnit.MILLISECONDS); + } + @Test public void asyncReceiveWithDefaultChannel() throws Exception { QueueChannel channel = new QueueChannel(); From e5f67740e4d3d08680a0fa3e7dda9971175c5d72 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Mon, 13 Sep 2010 14:51:45 -0400 Subject: [PATCH 2/7] INT-1254, fixed parent pom and security pom to pick up version of spring security from parent definition --- spring-integration-parent/pom.xml | 76 +++++++++++++---------------- spring-integration-security/pom.xml | 2 - 2 files changed, 35 insertions(+), 43 deletions(-) diff --git a/spring-integration-parent/pom.xml b/spring-integration-parent/pom.xml index 38b448079f..137d3c1f99 100644 --- a/spring-integration-parent/pom.xml +++ b/spring-integration-parent/pom.xml @@ -1,5 +1,6 @@ - + 4.0.0 org.springframework.integration spring-integration-parent @@ -23,7 +24,7 @@ 1.1 1.5.10 3.0.3.RELEASE - 2.0.5.RELEASE + 3.0.3.RELEASE 1.5.9 @@ -101,7 +102,8 @@ - + http://static.springframework.org/spring-integration/site/downloads/releases.html static.springframework.org @@ -109,14 +111,10 @@ - + org.aspectj @@ -213,13 +211,21 @@ spring-commons-serializer 1.0.0.M1 + + org.springframework.security + spring-security-core + ${org.springframework.security.version} + + + org.springframework.security + spring-security-config + ${org.springframework.security.version} + - + cglib cglib-nodep ${cglib.version} @@ -276,12 +282,10 @@ - + log4j log4j @@ -292,10 +296,8 @@ - + org.springframework.build.aws org.springframework.build.aws.maven 3.0.0.RELEASE @@ -376,12 +378,9 @@ - + com.springsource.bundlor com.springsource.bundlor.maven 1.0.0.RELEASE @@ -398,10 +397,8 @@ - + org.apache.maven.plugins maven-jar-plugin 2.2 @@ -416,11 +413,8 @@ - + org.apache.maven.plugins maven-project-info-reports-plugin 2.1 diff --git a/spring-integration-security/pom.xml b/spring-integration-security/pom.xml index 485eaa7f66..d76f292cff 100644 --- a/spring-integration-security/pom.xml +++ b/spring-integration-security/pom.xml @@ -27,7 +27,6 @@ org.springframework.security spring-security-core - 3.0.3.RELEASE org.springframework @@ -38,7 +37,6 @@ org.springframework.security spring-security-config - 3.0.3.RELEASE org.springframework From f7bca5e9ba33ab5bc4e02ed321032195e58b3831 Mon Sep 17 00:00:00 2001 From: Mark Fisher Date: Mon, 13 Sep 2010 15:05:10 -0400 Subject: [PATCH 3/7] INT-1440 the QoS properties are now only applied to JmsTemplate if it is not a referenced bean --- .../jms/AbstractJmsTemplateBasedAdapter.java | 26 +++++++------------ .../JmsOutboundChannelAdapterParserTests.java | 14 ++++++++++ 2 files changed, 24 insertions(+), 16 deletions(-) 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 2e1ede4682..f15440ce3e 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 @@ -191,21 +191,13 @@ public abstract class AbstractJmsTemplateBasedAdapter extends IntegrationObjectS && (this.destination != null || this.destinationName != null), "Either a 'jmsTemplate' or *both* 'connectionFactory' and" + " 'destination' (or 'destination-name') are required."); - this.jmsTemplate = this.createDefaultJmsTemplate(); + this.jmsTemplate = this.createJmsTemplate(); } - this.jmsTemplate.setExplicitQosEnabled(this.explicitQosEnabled); - this.jmsTemplate.setTimeToLive(this.timeToLive); - this.jmsTemplate.setPriority(this.priority); - this.jmsTemplate.setDeliveryMode(this.deliveryMode); - if (this.messageConverter != null) { - this.jmsTemplate.setMessageConverter(this.messageConverter); - } - //this.configureMessageConverter(this.jmsTemplate); this.initialized = true; } } - private JmsTemplate createDefaultJmsTemplate() { + private JmsTemplate createJmsTemplate() { JmsTemplate jmsTemplate = new JmsTemplate(); jmsTemplate.setConnectionFactory(this.connectionFactory); if (this.destination != null) { @@ -218,16 +210,18 @@ public abstract class AbstractJmsTemplateBasedAdapter extends IntegrationObjectS if (this.destinationResolver != null) { jmsTemplate.setDestinationResolver(this.destinationResolver); } + jmsTemplate.setExplicitQosEnabled(this.explicitQosEnabled); + jmsTemplate.setTimeToLive(this.timeToLive); + jmsTemplate.setPriority(this.priority); + jmsTemplate.setDeliveryMode(this.deliveryMode); + if (this.messageConverter != null) { + jmsTemplate.setMessageConverter(this.messageConverter); + } return jmsTemplate; } -// protected void configureMessageConverter(JmsTemplate jmsTemplate) { -// MessageConverter converter = jmsTemplate.getMessageConverter(); -// if (converter == null) { -// jmsTemplate.setMessageConverter(new SimpleMessageConverter()); -// } -// } protected boolean shouldExtractPayload() { return extractPayload; } + } diff --git a/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/JmsOutboundChannelAdapterParserTests.java b/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/JmsOutboundChannelAdapterParserTests.java index 912c02f04e..40be26a252 100644 --- a/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/JmsOutboundChannelAdapterParserTests.java +++ b/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/JmsOutboundChannelAdapterParserTests.java @@ -101,6 +101,20 @@ public class JmsOutboundChannelAdapterParserTests { assertEquals(context.getBean("template"), jmsTemplate); } + @Test + public void adapterWithJmsTemplateQos() { + ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext( + "jmsOutboundWithJmsTemplateQos.xml", this.getClass()); + EventDrivenConsumer endpoint = (EventDrivenConsumer) context.getBean("adapter"); + DirectFieldAccessor handlerAccessor = new DirectFieldAccessor(new DirectFieldAccessor(endpoint).getPropertyValue("handler")); + JmsTemplate jmsTemplate = (JmsTemplate) handlerAccessor.getPropertyValue("jmsTemplate"); + assertNotNull(jmsTemplate); + assertEquals(context.getBean("template"), jmsTemplate); + assertTrue(jmsTemplate.isExplicitQosEnabled()); + assertEquals(7, jmsTemplate.getPriority()); + assertEquals(12345, jmsTemplate.getTimeToLive()); + } + @Test public void adapterWithMessageConverter() { ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext( From 4db74cd558eb90760cda3181f2666acf9812f651 Mon Sep 17 00:00:00 2001 From: Mark Fisher Date: Mon, 13 Sep 2010 15:06:25 -0400 Subject: [PATCH 4/7] INT-1440 the QoS properties are now only applied to JmsTemplate if it is not a referenced bean --- .../config/jmsOutboundWithJmsTemplateQos.xml | 32 +++++++++++++++++++ 1 file changed, 32 insertions(+) create mode 100644 spring-integration-jms/src/test/java/org/springframework/integration/jms/config/jmsOutboundWithJmsTemplateQos.xml diff --git a/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/jmsOutboundWithJmsTemplateQos.xml b/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/jmsOutboundWithJmsTemplateQos.xml new file mode 100644 index 0000000000..77bb36941f --- /dev/null +++ b/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/jmsOutboundWithJmsTemplateQos.xml @@ -0,0 +1,32 @@ + + + + + + + + + + + + + + + + + + + + + + + From 62f80d822b1ac18a5ecebcd22b594be72c51d80c Mon Sep 17 00:00:00 2001 From: Mark Fisher Date: Mon, 13 Sep 2010 15:18:42 -0400 Subject: [PATCH 5/7] removed System.out.println calls from tests --- .../ExceptionHandlingSiConsumerTests.java | 45 +++++++++++-------- .../jms/config/JmsMessageHistoryTests.java | 1 - 2 files changed, 27 insertions(+), 19 deletions(-) diff --git a/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/ExceptionHandlingSiConsumerTests.java b/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/ExceptionHandlingSiConsumerTests.java index ee09b2a3d6..4eddd20e73 100644 --- a/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/ExceptionHandlingSiConsumerTests.java +++ b/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/ExceptionHandlingSiConsumerTests.java @@ -13,8 +13,12 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.integration.jms.config; +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotNull; + import javax.jms.ConnectionFactory; import javax.jms.Destination; import javax.jms.JMSException; @@ -22,9 +26,8 @@ import javax.jms.Message; import javax.jms.Session; import javax.jms.TextMessage; -import org.junit.Assert; -import org.junit.Ignore; import org.junit.Test; + import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.support.ClassPathXmlApplicationContext; import org.springframework.integration.mapping.InboundMessageMapper; @@ -55,10 +58,10 @@ public class ExceptionHandlingSiConsumerTests { } }); Message message = jmsTemplate.receive(reply); - System.out.println(message); - Assert.assertNotNull(message); + assertNotNull(message); applicationContext.close(); } + @Test public void nonSiProducer_siConsumer_sync_withReturnNoException() throws Exception { ActiveMqTestUtils.prepare(); @@ -76,7 +79,8 @@ public class ExceptionHandlingSiConsumerTests { } }); Message message = jmsTemplate.receive(reply); - Assert.assertNotNull(message); + assertNotNull(message); + assertEquals("echoWithException", ((TextMessage) message).getText()); applicationContext.close(); } @@ -86,35 +90,40 @@ public class ExceptionHandlingSiConsumerTests { final ConfigurableApplicationContext applicationContext = new ClassPathXmlApplicationContext("Exception-nonSiProducer-siConsumer.xml", ExceptionHandlingSiConsumerTests.class); SampleGateway gateway = applicationContext.getBean("sampleGateway", SampleGateway.class); String reply = gateway.echo("echoWithExceptionChannel"); - System.out.println("Reply: " + reply); + assertEquals("echoWithException", reply); applicationContext.close(); } -// - public static class SampleService{ - public String echoWithException(String value){ + + + public static class SampleService { + + public String echoWithException(String value) { throw new SampleException("echoWithException"); } + public String echo(String value){ return value; } } - - + + @SuppressWarnings("serial") - public static class SampleException extends RuntimeException{ + public static class SampleException extends RuntimeException { public SampleException(String message){ super(message); } } - - public static interface SampleGateway{ + + + public static interface SampleGateway { public String echo(String value); } - - public static class SampleErrorMessageMapper implements InboundMessageMapper{ - public org.springframework.integration.Message toMessage( - Throwable t) throws Exception { + + + public static class SampleErrorMessageMapper implements InboundMessageMapper { + public org.springframework.integration.Message toMessage(Throwable t) throws Exception { return MessageBuilder.withPayload(t.getCause().getMessage()).build(); } } + } 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 index b1733a89ca..4e78ba5683 100644 --- 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 @@ -58,7 +58,6 @@ public class JmsMessageHistoryTests { assertEquals("jms:inbound-channel-adapter", event1.getProperty(MessageHistory.TYPE_PROPERTY)); assertEquals("sampleJmsInboundAdapter", event1.getProperty(MessageHistory.NAME_PROPERTY)); Properties event2 = historyIterator.next(); - System.out.println(event2); assertEquals("channel", event2.getProperty(MessageHistory.TYPE_PROPERTY)); assertEquals("jmsInputChannel", event2.getProperty(MessageHistory.NAME_PROPERTY)); } From f964f74f4e0653fecc8e260fec170db3088da880 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Mon, 13 Sep 2010 15:57:57 -0400 Subject: [PATCH 6/7] INT-1448, INT-1442 accepted changes by Iwein, added small refactoring to AbstractMessageSplitter, added support for 'requires-reply' attribute, added test cases' --- .../config/SplitterFactoryBean.java | 8 +++++++ .../splitter/AbstractMessageSplitter.java | 12 ++++++----- .../config/xml/spring-integration-2.0.xsd | 9 ++++++++ .../router/config/SplitterParserTests.java | 15 +++++++++++++ .../router/config/splitterParserTests.xml | 5 +++++ .../splitter/DefaultSplitterTests.java | 21 ++++++++++--------- 6 files changed, 55 insertions(+), 15 deletions(-) diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/SplitterFactoryBean.java b/spring-integration-core/src/main/java/org/springframework/integration/config/SplitterFactoryBean.java index 706dd8095a..dbae530ef5 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/SplitterFactoryBean.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/SplitterFactoryBean.java @@ -31,6 +31,7 @@ import org.springframework.util.StringUtils; public class SplitterFactoryBean extends AbstractMessageHandlerFactoryBean { private volatile Long sendTimeout; + private volatile boolean requiresReply; public void setSendTimeout(Long sendTimeout) { this.sendTimeout = sendTimeout; @@ -64,7 +65,14 @@ public class SplitterFactoryBean extends AbstractMessageHandlerFactoryBean { if (this.sendTimeout != null) { splitter.setSendTimeout(sendTimeout); } + splitter.setRequiresReply(requiresReply); return splitter; } + public boolean isRequiresReply() { + return requiresReply; + } + public void setRequiresReply(boolean requiresReply) { + this.requiresReply = requiresReply; + } } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/splitter/AbstractMessageSplitter.java b/spring-integration-core/src/main/java/org/springframework/integration/splitter/AbstractMessageSplitter.java index 161c6c8b47..a53added58 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/splitter/AbstractMessageSplitter.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/splitter/AbstractMessageSplitter.java @@ -20,7 +20,10 @@ import org.springframework.integration.Message; import org.springframework.integration.MessageHeaders; import org.springframework.integration.handler.AbstractReplyProducingMessageHandler; import org.springframework.integration.support.MessageBuilder; +import org.springframework.util.CollectionUtils; +import org.springframework.util.ObjectUtils; +import java.lang.reflect.Array; import java.util.*; /** @@ -38,7 +41,10 @@ public abstract class AbstractMessageSplitter extends AbstractReplyProducingMess @SuppressWarnings("unchecked") protected final Object handleRequestMessage(Message message) { Object result = this.splitMessage(message); - if (result == null) { + // return null if 'null', empty Collection or empty Array + if ( result == null || + (result instanceof Collection && CollectionUtils.isEmpty((Collection)result)) || + (result.getClass().isArray() && ObjectUtils.isEmpty((Object[]) result)) ) { return null; } MessageHeaders headers = message.getHeaders(); @@ -59,10 +65,6 @@ public abstract class AbstractMessageSplitter extends AbstractReplyProducingMess List> messageBuilders = new ArrayList>(); if (result instanceof Collection) { Collection items = (Collection) result; - //TODO put this return statement in a more obvious place - if(items.isEmpty()){ - return null; - } int sequenceNumber = 0; int sequenceSize = items.size(); for (Object item : items) { diff --git a/spring-integration-core/src/main/resources/org/springframework/integration/config/xml/spring-integration-2.0.xsd b/spring-integration-core/src/main/resources/org/springframework/integration/config/xml/spring-integration-2.0.xsd index 760f660fc3..7066cebf9b 100644 --- a/spring-integration-core/src/main/resources/org/springframework/integration/config/xml/spring-integration-2.0.xsd +++ b/spring-integration-core/src/main/resources/org/springframework/integration/config/xml/spring-integration-2.0.xsd @@ -1948,6 +1948,15 @@ Name of the header whose value to use. + + + + Specify whether the splitter method must return a non-null value. This value will be + FALSE by default, but if set to TRUE, a MessageHandlingException will be thrown when + the underlying service method (or expression) returns a NULL value. + + + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/router/config/SplitterParserTests.java b/spring-integration-core/src/test/java/org/springframework/integration/router/config/SplitterParserTests.java index 3a795055c8..8e13d1a5d6 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/router/config/SplitterParserTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/router/config/SplitterParserTests.java @@ -19,13 +19,19 @@ package org.springframework.integration.router.config; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertNull; +import java.util.Collections; + import org.junit.Test; import org.springframework.context.support.ClassPathXmlApplicationContext; import org.springframework.integration.Message; import org.springframework.integration.MessageChannel; +import org.springframework.integration.MessageHandlingException; +import org.springframework.integration.channel.DirectChannel; import org.springframework.integration.core.PollableChannel; +import org.springframework.integration.endpoint.EventDrivenConsumer; import org.springframework.integration.message.GenericMessage; +import org.springframework.integration.support.MessageBuilder; /** * @author Mark Fisher @@ -88,5 +94,14 @@ public class SplitterParserTests { assertEquals("test", result4.getPayload()); assertNull(output.receive(0)); } + + @Test(expected=MessageHandlingException.class) + public void splitterParserTestWithRequiresReply() { + ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext( + "splitterParserTests.xml", this.getClass()); + context.start(); + DirectChannel inputChannel = context.getBean("requiresReplyInput", DirectChannel.class); + inputChannel.send(MessageBuilder.withPayload(Collections.emptyList()).build()); + } } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/router/config/splitterParserTests.xml b/spring-integration-core/src/test/java/org/springframework/integration/router/config/splitterParserTests.xml index 190a41014b..aa96b3df2b 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/router/config/splitterParserTests.xml +++ b/spring-integration-core/src/test/java/org/springframework/integration/router/config/splitterParserTests.xml @@ -26,6 +26,11 @@ ref="splitterImpl" input-channel="splitterImplementationInput" output-channel="output"/> + + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/splitter/DefaultSplitterTests.java b/spring-integration-core/src/test/java/org/springframework/integration/splitter/DefaultSplitterTests.java index e2326ec28a..bcfafb3f93 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/splitter/DefaultSplitterTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/splitter/DefaultSplitterTests.java @@ -16,6 +16,17 @@ package org.springframework.integration.splitter; +import static junit.framework.Assert.assertEquals; +import static org.hamcrest.Matchers.is; +import static org.hamcrest.Matchers.nullValue; +import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertThat; +import static org.junit.Assert.assertTrue; + +import java.util.Arrays; +import java.util.Collections; +import java.util.List; + import org.junit.Test; import org.springframework.integration.Message; import org.springframework.integration.channel.DirectChannel; @@ -23,15 +34,6 @@ import org.springframework.integration.channel.QueueChannel; import org.springframework.integration.endpoint.EventDrivenConsumer; import org.springframework.integration.support.MessageBuilder; -import java.util.Arrays; -import java.util.Collections; -import java.util.List; - -import static junit.framework.Assert.assertEquals; -import static org.hamcrest.Matchers.is; -import static org.hamcrest.Matchers.nullValue; -import static org.junit.Assert.*; - /** * @author Mark Fisher * @author Iwein Fuld @@ -104,5 +106,4 @@ public class DefaultSplitterTests { Message output = replyChannel.receive(15); assertThat(output, is(nullValue())); } - } From 3109a781058071ca95d277bb22b6abfbb66a8db3 Mon Sep 17 00:00:00 2001 From: Mark Fisher Date: Mon, 13 Sep 2010 16:01:59 -0400 Subject: [PATCH 7/7] INT-1447 raising an error if any JmsTemplate attributes are set when an explicit 'jms-template' reference is also specified --- .../jms/config/JmsAdapterParserUtils.java | 13 +++++++++++++ .../jms/config/JmsInboundChannelAdapterParser.java | 10 +--------- .../jms/config/JmsOutboundChannelAdapterParser.java | 11 +++-------- 3 files changed, 17 insertions(+), 17 deletions(-) diff --git a/spring-integration-jms/src/main/java/org/springframework/integration/jms/config/JmsAdapterParserUtils.java b/spring-integration-jms/src/main/java/org/springframework/integration/jms/config/JmsAdapterParserUtils.java index c0a8d3d365..fce101b8ee 100644 --- a/spring-integration-jms/src/main/java/org/springframework/integration/jms/config/JmsAdapterParserUtils.java +++ b/spring-integration-jms/src/main/java/org/springframework/integration/jms/config/JmsAdapterParserUtils.java @@ -52,6 +52,10 @@ abstract class JmsAdapterParserUtils { static final String HEADER_MAPPER_PROPERTY = "headerMapper"; + private static final String[] JMS_TEMPLATE_ATTRIBUTES = { "destination", "destination-name", + "connection-factory", "message-converter", "time-to-live", "priority", "delivery-persistent", "explicit-qos-enabled" }; + + /* * The following constants match those of javax.jms.Session. * They are duplicated here to avoid a dependency in tooling. @@ -102,4 +106,13 @@ abstract class JmsAdapterParserUtils { } } + static void verifyNoJmsTemplateAttributes(Element element, ParserContext parserContext) { + for (String attributeName : JMS_TEMPLATE_ATTRIBUTES) { + if (element.hasAttribute(attributeName)) { + parserContext.getReaderContext().error("When providing a 'jms-template' reference, the '" + + attributeName + "' attribute is not allowed", parserContext.extractSource(element)); + } + } + } + } 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 b952fbfa06..3ba23c5e6b 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 @@ -60,15 +60,7 @@ public class JmsInboundChannelAdapterParser extends AbstractPollingInboundChanne boolean hasDestinationRef = StringUtils.hasText(destination); boolean hasDestinationName = StringUtils.hasText(destinationName); if (StringUtils.hasText(jmsTemplate)) { - if (element.hasAttribute(JmsAdapterParserUtils.CONNECTION_FACTORY_ATTRIBUTE) || - hasDestinationRef || hasDestinationName) { - parserContext.getReaderContext().error( - "When providing '" + JmsAdapterParserUtils.JMS_TEMPLATE_ATTRIBUTE + - "', none of '" + JmsAdapterParserUtils.CONNECTION_FACTORY_ATTRIBUTE + - "', '" + JmsAdapterParserUtils.DESTINATION_ATTRIBUTE + "', or '" + - JmsAdapterParserUtils.DESTINATION_NAME_ATTRIBUTE + "' are allowed.", - source); - } + JmsAdapterParserUtils.verifyNoJmsTemplateAttributes(element, parserContext); builder.addConstructorArgReference(jmsTemplate); } else if (hasDestinationRef || hasDestinationName) { diff --git a/spring-integration-jms/src/main/java/org/springframework/integration/jms/config/JmsOutboundChannelAdapterParser.java b/spring-integration-jms/src/main/java/org/springframework/integration/jms/config/JmsOutboundChannelAdapterParser.java index 16624af905..e1c221ac53 100644 --- a/spring-integration-jms/src/main/java/org/springframework/integration/jms/config/JmsOutboundChannelAdapterParser.java +++ b/spring-integration-jms/src/main/java/org/springframework/integration/jms/config/JmsOutboundChannelAdapterParser.java @@ -18,7 +18,6 @@ package org.springframework.integration.jms.config; import org.w3c.dom.Element; -import org.springframework.beans.factory.BeanCreationException; import org.springframework.beans.factory.support.AbstractBeanDefinition; import org.springframework.beans.factory.support.BeanDefinitionBuilder; import org.springframework.beans.factory.xml.ParserContext; @@ -44,11 +43,7 @@ public class JmsOutboundChannelAdapterParser extends AbstractOutboundChannelAdap boolean hasDestinationRef = StringUtils.hasText(destination); boolean hasDestinationName = StringUtils.hasText(destinationName); if (StringUtils.hasText(jmsTemplate)) { - if (element.hasAttribute(JmsAdapterParserUtils.CONNECTION_FACTORY_ATTRIBUTE) || - hasDestinationRef || hasDestinationName) { - throw new BeanCreationException("When providing a 'jms-template' reference, none of " + - "'connection-factory', 'destination', or 'destination-name' should be provided."); - } + JmsAdapterParserUtils.verifyNoJmsTemplateAttributes(element, parserContext); builder.addConstructorArgReference(jmsTemplate); } else if (hasDestinationRef ^ hasDestinationName) { @@ -64,8 +59,8 @@ public class JmsOutboundChannelAdapterParser extends AbstractOutboundChannelAdap } } else { - throw new BeanCreationException("Either a 'jms-template' reference " + - "or one of 'destination' or 'destination-name' must be provided."); + parserContext.getReaderContext().error("Either a 'jms-template' reference " + + "or one of 'destination' or 'destination-name' must be provided.", parserContext.extractSource(element)); } if (StringUtils.hasText(headerMapper)) { builder.addPropertyReference(JmsAdapterParserUtils.HEADER_MAPPER_PROPERTY, headerMapper);