();
+ ReflectionUtils.doWithMethods(bean.getClass(), new ReflectionUtils.MethodCallback() {
+ @Override
+ public void doWith(Method method) throws IllegalArgumentException, IllegalAccessException {
+ Annotation annotation = AnnotationUtils.getAnnotation(method, CorrelationStrategy.class);
+ if (annotation != null) {
+ reference.set(new MethodInvokingCorrelationStrategy(bean, method));
+ }
+ }
+ });
+ return reference.get();
+ }
}
diff --git a/spring-integration-core/src/main/java/org/springframework/integration/dispatcher/AbstractDispatcher.java b/spring-integration-core/src/main/java/org/springframework/integration/dispatcher/AbstractDispatcher.java
index 6eb36a5d85..d37db72fa2 100644
--- a/spring-integration-core/src/main/java/org/springframework/integration/dispatcher/AbstractDispatcher.java
+++ b/spring-integration-core/src/main/java/org/springframework/integration/dispatcher/AbstractDispatcher.java
@@ -23,8 +23,6 @@ import org.apache.commons.logging.LogFactory;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageDeliveryException;
-import org.springframework.integration.support.DefaultMessageBuilderFactory;
-import org.springframework.integration.support.MessageBuilderFactory;
import org.springframework.messaging.MessageHandler;
import org.springframework.messaging.MessagingException;
import org.springframework.util.Assert;
@@ -55,7 +53,6 @@ public abstract class AbstractDispatcher implements MessageDispatcher {
private volatile MessageHandler theOneHandler;
- private volatile MessageBuilderFactory messageBuilderFactory = new DefaultMessageBuilderFactory();
/**
* Set the maximum subscribers allowed by this dispatcher.
* @param maxSubscribers The maximum number of subscribers allowed.
@@ -74,15 +71,6 @@ public abstract class AbstractDispatcher implements MessageDispatcher {
return handlers.asUnmodifiableSet();
}
- protected MessageBuilderFactory getMessageBuilderFactory() {
- return messageBuilderFactory;
- }
-
- public void setMessageBuilderFactory(MessageBuilderFactory messageBuilderFactory) {
- Assert.notNull(messageBuilderFactory, "'messageBuilderFactory' cannot be null");
- this.messageBuilderFactory = messageBuilderFactory;
- }
-
/**
* Add the handler to the internal Set.
*
diff --git a/spring-integration-core/src/main/java/org/springframework/integration/dispatcher/BroadcastingDispatcher.java b/spring-integration-core/src/main/java/org/springframework/integration/dispatcher/BroadcastingDispatcher.java
index e1a8f01757..0476ec1076 100644
--- a/spring-integration-core/src/main/java/org/springframework/integration/dispatcher/BroadcastingDispatcher.java
+++ b/spring-integration-core/src/main/java/org/springframework/integration/dispatcher/BroadcastingDispatcher.java
@@ -19,7 +19,13 @@ package org.springframework.integration.dispatcher;
import java.util.Collection;
import java.util.concurrent.Executor;
+import org.springframework.beans.BeansException;
+import org.springframework.beans.factory.BeanFactory;
+import org.springframework.beans.factory.BeanFactoryAware;
import org.springframework.integration.MessageDispatchingException;
+import org.springframework.integration.context.IntegrationContextUtils;
+import org.springframework.integration.support.DefaultMessageBuilderFactory;
+import org.springframework.integration.support.MessageBuilderFactory;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageHandler;
import org.springframework.messaging.MessagingException;
@@ -39,7 +45,7 @@ import org.springframework.messaging.MessagingException;
* @author Gary Russell
* @author Oleg Zhurakousky
*/
-public class BroadcastingDispatcher extends AbstractDispatcher {
+public class BroadcastingDispatcher extends AbstractDispatcher implements BeanFactoryAware {
private final boolean requireSubscribers;
@@ -51,6 +57,9 @@ public class BroadcastingDispatcher extends AbstractDispatcher {
private volatile int minSubscribers;
+ private volatile MessageBuilderFactory messageBuilderFactory = new DefaultMessageBuilderFactory();
+
+
public BroadcastingDispatcher() {
this(null, false);
}
@@ -103,6 +112,11 @@ public class BroadcastingDispatcher extends AbstractDispatcher {
this.minSubscribers = minSubscribers;
}
+ @Override
+ public void setBeanFactory(BeanFactory beanFactory) throws BeansException {
+ this.messageBuilderFactory = IntegrationContextUtils.getMessageBuilderFactory(beanFactory);
+ }
+
@Override
public boolean dispatch(Message> message) {
int dispatched = 0;
@@ -113,7 +127,7 @@ public class BroadcastingDispatcher extends AbstractDispatcher {
}
int sequenceSize = handlers.size();
for (final MessageHandler handler : handlers) {
- final Message> messageToSend = (!this.applySequence) ? message : this.getMessageBuilderFactory().fromMessage(message)
+ final Message> messageToSend = (!this.applySequence) ? message : this.messageBuilderFactory.fromMessage(message)
.pushSequenceDetails(message.getHeaders().getId(), sequenceNumber++, sequenceSize).build();
if (this.executor != null) {
this.executor.execute(new Runnable() {
diff --git a/spring-integration-core/src/main/java/org/springframework/integration/gateway/MessagingGatewaySupport.java b/spring-integration-core/src/main/java/org/springframework/integration/gateway/MessagingGatewaySupport.java
index da4dbf26e3..c2aace6fde 100644
--- a/spring-integration-core/src/main/java/org/springframework/integration/gateway/MessagingGatewaySupport.java
+++ b/spring-integration-core/src/main/java/org/springframework/integration/gateway/MessagingGatewaySupport.java
@@ -180,6 +180,7 @@ public abstract class MessagingGatewaySupport extends AbstractEndpoint implement
if (this.requestMapper instanceof DefaultRequestMapper) {
((DefaultRequestMapper) this.requestMapper).setMessageBuilderFactory(this.getMessageBuilderFactory());
}
+ this.messageConverter.setBeanFactory(this.getBeanFactory());
}
this.initialized = true;
}
diff --git a/spring-integration-core/src/main/java/org/springframework/integration/support/converter/SimpleMessageConverter.java b/spring-integration-core/src/main/java/org/springframework/integration/support/converter/SimpleMessageConverter.java
index ee16478efd..f5f7b97aa2 100644
--- a/spring-integration-core/src/main/java/org/springframework/integration/support/converter/SimpleMessageConverter.java
+++ b/spring-integration-core/src/main/java/org/springframework/integration/support/converter/SimpleMessageConverter.java
@@ -16,6 +16,10 @@
package org.springframework.integration.support.converter;
+import org.springframework.beans.BeansException;
+import org.springframework.beans.factory.BeanFactory;
+import org.springframework.beans.factory.BeanFactoryAware;
+import org.springframework.integration.context.IntegrationContextUtils;
import org.springframework.integration.mapping.InboundMessageMapper;
import org.springframework.integration.mapping.OutboundMessageMapper;
import org.springframework.integration.support.DefaultMessageBuilderFactory;
@@ -31,7 +35,7 @@ import org.springframework.messaging.converter.MessageConverter;
* @since 2.0
*/
@SuppressWarnings({"unchecked", "rawtypes"})
-public class SimpleMessageConverter implements MessageConverter {
+public class SimpleMessageConverter implements MessageConverter, BeanFactoryAware {
private volatile InboundMessageMapper inboundMessageMapper;
@@ -67,8 +71,9 @@ public class SimpleMessageConverter implements MessageConverter {
this.outboundMessageMapper = (outboundMessageMapper != null) ? outboundMessageMapper : new DefaultOutboundMessageMapper();
}
- public final void setMessageBuilderFactory(MessageBuilderFactory messageBuilderFactory) {
- this.messageBuilderFactory = messageBuilderFactory;
+ @Override
+ public void setBeanFactory(BeanFactory beanFactory) throws BeansException {
+ this.messageBuilderFactory = IntegrationContextUtils.getMessageBuilderFactory(beanFactory);
}
@Override
diff --git a/spring-integration-core/src/main/java/org/springframework/integration/support/json/AbstractJacksonJsonMessageParser.java b/spring-integration-core/src/main/java/org/springframework/integration/support/json/AbstractJacksonJsonMessageParser.java
index ec63763c3c..c27774e1a8 100644
--- a/spring-integration-core/src/main/java/org/springframework/integration/support/json/AbstractJacksonJsonMessageParser.java
+++ b/spring-integration-core/src/main/java/org/springframework/integration/support/json/AbstractJacksonJsonMessageParser.java
@@ -18,10 +18,13 @@ package org.springframework.integration.support.json;
import java.lang.reflect.Type;
+import org.springframework.beans.BeansException;
+import org.springframework.beans.factory.BeanFactory;
+import org.springframework.beans.factory.BeanFactoryAware;
+import org.springframework.integration.context.IntegrationContextUtils;
import org.springframework.integration.support.DefaultMessageBuilderFactory;
import org.springframework.integration.support.MessageBuilderFactory;
import org.springframework.messaging.Message;
-import org.springframework.util.Assert;
/**
* Base {@link JsonInboundMessageMapper.JsonMessageParser} implementation for Jackson processors.
@@ -30,7 +33,8 @@ import org.springframework.util.Assert;
* @since 3.0
*
*/
-abstract class AbstractJacksonJsonMessageParser implements JsonInboundMessageMapper.JsonMessageParser
{
+abstract class AbstractJacksonJsonMessageParser
implements JsonInboundMessageMapper.JsonMessageParser
,
+ BeanFactoryAware {
private final JsonObjectMapper, P> objectMapper;
@@ -42,9 +46,9 @@ abstract class AbstractJacksonJsonMessageParser
implements JsonInboundMessage
this.objectMapper = objectMapper;
}
- public void setMessageBuilderFactory(MessageBuilderFactory messageBuilderFactory) {
- Assert.notNull(messageBuilderFactory, "'messageBuilderFactory' cannot be null");
- this.messageBuilderFactory = messageBuilderFactory;
+ @Override
+ public void setBeanFactory(BeanFactory beanFactory) throws BeansException {
+ this.messageBuilderFactory = IntegrationContextUtils.getMessageBuilderFactory(beanFactory);
}
protected MessageBuilderFactory getMessageBuilderFactory() {
diff --git a/spring-integration-core/src/main/java/org/springframework/integration/util/AbstractExpressionEvaluator.java b/spring-integration-core/src/main/java/org/springframework/integration/util/AbstractExpressionEvaluator.java
index fbc1795c85..330d3c0b95 100644
--- a/spring-integration-core/src/main/java/org/springframework/integration/util/AbstractExpressionEvaluator.java
+++ b/spring-integration-core/src/main/java/org/springframework/integration/util/AbstractExpressionEvaluator.java
@@ -32,7 +32,6 @@ import org.springframework.integration.support.DefaultMessageBuilderFactory;
import org.springframework.integration.support.MessageBuilderFactory;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageHandlingException;
-import org.springframework.util.Assert;
/**
* @author Mark Fisher
@@ -60,6 +59,7 @@ public abstract class AbstractExpressionEvaluator implements BeanFactoryAware, I
/**
* Specify a BeanFactory in order to enable resolution via @beanName in the expression.
*/
+ @Override
public void setBeanFactory(final BeanFactory beanFactory) {
if (beanFactory != null) {
this.beanFactory = beanFactory;
@@ -67,6 +67,7 @@ public abstract class AbstractExpressionEvaluator implements BeanFactoryAware, I
if (this.evaluationContext != null && this.evaluationContext.getBeanResolver() == null) {
this.evaluationContext.setBeanResolver(new BeanFactoryResolver(beanFactory));
}
+ this.messageBuilderFactory = IntegrationContextUtils.getMessageBuilderFactory(beanFactory);
}
}
@@ -76,11 +77,6 @@ public abstract class AbstractExpressionEvaluator implements BeanFactoryAware, I
}
}
- public void setMessageBuilderFactory(MessageBuilderFactory messageBuilderFactory) {
- Assert.notNull(messageBuilderFactory, "'messageBuilderFactory' cannot be null");
- this.messageBuilderFactory = messageBuilderFactory;
- }
-
protected MessageBuilderFactory getMessageBuilderFactory() {
if (this.messageBuilderFactory == null) {
this.messageBuilderFactory = new DefaultMessageBuilderFactory();
@@ -92,7 +88,7 @@ public abstract class AbstractExpressionEvaluator implements BeanFactoryAware, I
public void afterPropertiesSet() throws Exception {
getEvaluationContext();
if (this.messageBuilderFactory == null) {
- this.messageBuilderFactory = IntegrationContextUtils.getMessageBuilderFactory(beanFactory);
+ this.messageBuilderFactory = IntegrationContextUtils.getMessageBuilderFactory(this.beanFactory);
}
}
diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/AggregatorParserTests.java b/spring-integration-core/src/test/java/org/springframework/integration/config/AggregatorParserTests.java
index e488aed7d8..e25aad5a70 100644
--- a/spring-integration-core/src/test/java/org/springframework/integration/config/AggregatorParserTests.java
+++ b/spring-integration-core/src/test/java/org/springframework/integration/config/AggregatorParserTests.java
@@ -1,5 +1,5 @@
/*
- * Copyright 2002-2013 the original author or authors.
+ * Copyright 2002-2014 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.
@@ -37,7 +37,6 @@ import org.springframework.beans.factory.BeanCreationException;
import org.springframework.beans.factory.parsing.BeanDefinitionParsingException;
import org.springframework.context.ApplicationContext;
import org.springframework.context.support.ClassPathXmlApplicationContext;
-import org.springframework.messaging.MessageHandlingException;
import org.springframework.integration.MessageRejectedException;
import org.springframework.integration.aggregator.AggregatingMessageHandler;
import org.springframework.integration.aggregator.CorrelationStrategy;
@@ -46,6 +45,7 @@ import org.springframework.integration.aggregator.ExpressionEvaluatingReleaseStr
import org.springframework.integration.aggregator.MethodInvokingMessageGroupProcessor;
import org.springframework.integration.aggregator.MethodInvokingReleaseStrategy;
import org.springframework.integration.aggregator.ReleaseStrategy;
+import org.springframework.integration.context.IntegrationContextUtils;
import org.springframework.integration.endpoint.EventDrivenConsumer;
import org.springframework.integration.support.MessageBuilder;
import org.springframework.integration.test.util.TestUtils;
@@ -53,6 +53,7 @@ import org.springframework.messaging.Message;
import org.springframework.messaging.MessageChannel;
import org.springframework.messaging.MessageDeliveryException;
import org.springframework.messaging.MessageHandler;
+import org.springframework.messaging.MessageHandlingException;
import org.springframework.messaging.PollableChannel;
import org.springframework.messaging.SubscribableChannel;
@@ -63,6 +64,7 @@ import org.springframework.messaging.SubscribableChannel;
* @author Oleg Zhurakousky
* @author Artem Bilan
* @author Gunnar Hillert
+ * @author Gary Russell
*/
public class AggregatorParserTests {
@@ -88,6 +90,10 @@ public class AggregatorParserTests {
.size());
Message> aggregatedMessage = aggregatorBean.getAggregatedMessages().get("id1");
assertEquals("The aggregated message payload is not correct", "123456789", aggregatedMessage.getPayload());
+ Object mbf = context.getBean(IntegrationContextUtils.INTEGRATION_MESSAGE_BUILDER_FACTORY_BEAN_NAME);
+ Object handler = context.getBean("aggregatorWithReference.handler");
+ assertSame(mbf, TestUtils.getPropertyValue(handler, "outputProcessor.messageBuilderFactory"));
+ assertSame(mbf, TestUtils.getPropertyValue(handler, "outputProcessor.processor.messageBuilderFactory"));
}
@Test
@@ -96,6 +102,7 @@ public class AggregatorParserTests {
SubscribableChannel outputChannel = (SubscribableChannel) context.getBean("aggregatorWithExpressionsOutput");
final AtomicReference> aggregatedMessage = new AtomicReference>();
outputChannel.subscribe(new MessageHandler() {
+ @Override
public void handleMessage(Message> message) throws MessageRejectedException, MessageHandlingException,
MessageDeliveryException {
aggregatedMessage.set(message);
@@ -110,6 +117,10 @@ public class AggregatorParserTests {
}
assertEquals("The aggregated message payload is not correct", "[123]", aggregatedMessage.get().getPayload()
.toString());
+ Object mbf = context.getBean(IntegrationContextUtils.INTEGRATION_MESSAGE_BUILDER_FACTORY_BEAN_NAME);
+ Object handler = context.getBean("aggregatorWithExpressions.handler");
+ assertSame(mbf, TestUtils.getPropertyValue(handler, "outputProcessor.messageBuilderFactory"));
+ assertSame(mbf, TestUtils.getPropertyValue(handler, "outputProcessor.processor.messageBuilderFactory"));
}
@Test
@@ -158,6 +169,10 @@ public class AggregatorParserTests {
PollableChannel outputChannel = (PollableChannel) context.getBean("outputChannel");
Message> response = outputChannel.receive(10);
Assert.assertEquals(6l, response.getPayload());
+ Object mbf = context.getBean(IntegrationContextUtils.INTEGRATION_MESSAGE_BUILDER_FACTORY_BEAN_NAME);
+ Object handler = context.getBean("aggregatorWithReferenceAndMethod.handler");
+ assertSame(mbf, TestUtils.getPropertyValue(handler, "outputProcessor.messageBuilderFactory"));
+ assertSame(mbf, TestUtils.getPropertyValue(handler, "outputProcessor.processor.messageBuilderFactory"));
}
@Test(expected = BeanCreationException.class)
@@ -230,6 +245,8 @@ public class AggregatorParserTests {
EventDrivenConsumer aggregatorConsumer = (EventDrivenConsumer) context.getBean("aggregatorWithExpressionsAndPojoAggregator");
AggregatingMessageHandler aggregatingMessageHandler = (AggregatingMessageHandler) TestUtils.getPropertyValue(aggregatorConsumer, "handler");
MethodInvokingMessageGroupProcessor messageGroupProcessor = (MethodInvokingMessageGroupProcessor) TestUtils.getPropertyValue(aggregatingMessageHandler, "outputProcessor");
+ Object mbf = context.getBean(IntegrationContextUtils.INTEGRATION_MESSAGE_BUILDER_FACTORY_BEAN_NAME);
+ assertSame(mbf, TestUtils.getPropertyValue(messageGroupProcessor, "messageBuilderFactory"));
Object messageGroupProcessorTargetObject = TestUtils.getPropertyValue(messageGroupProcessor, "processor.delegate.targetObject");
assertSame(context.getBean("aggregatorBean"), messageGroupProcessorTargetObject);
ReleaseStrategy releaseStrategy = (ReleaseStrategy) TestUtils.getPropertyValue(aggregatingMessageHandler, "releaseStrategy");
diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/PublishSubscribeChannelParserTests.java b/spring-integration-core/src/test/java/org/springframework/integration/config/PublishSubscribeChannelParserTests.java
index 14a7d33388..bfc72e5dd1 100644
--- a/spring-integration-core/src/test/java/org/springframework/integration/config/PublishSubscribeChannelParserTests.java
+++ b/spring-integration-core/src/test/java/org/springframework/integration/config/PublishSubscribeChannelParserTests.java
@@ -1,5 +1,5 @@
/*
- * Copyright 2002-2010 the original author or authors.
+ * Copyright 2002-2014 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.
@@ -20,6 +20,7 @@ import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.assertNull;
+import static org.junit.Assert.assertSame;
import static org.junit.Assert.assertTrue;
import java.util.concurrent.Executor;
@@ -29,12 +30,14 @@ import org.junit.Test;
import org.springframework.beans.DirectFieldAccessor;
import org.springframework.context.support.ClassPathXmlApplicationContext;
import org.springframework.integration.channel.PublishSubscribeChannel;
+import org.springframework.integration.context.IntegrationContextUtils;
import org.springframework.integration.dispatcher.BroadcastingDispatcher;
import org.springframework.integration.util.ErrorHandlingTaskExecutor;
import org.springframework.util.ErrorHandler;
/**
* @author Mark Fisher
+ * @author Gary Russell
*/
public class PublishSubscribeChannelParserTests {
@@ -51,6 +54,9 @@ public class PublishSubscribeChannelParserTests {
assertNull(dispatcherAccessor.getPropertyValue("executor"));
assertFalse((Boolean) dispatcherAccessor.getPropertyValue("ignoreFailures"));
assertFalse((Boolean) dispatcherAccessor.getPropertyValue("applySequence"));
+ Object mbf = context.getBean(IntegrationContextUtils.INTEGRATION_MESSAGE_BUILDER_FACTORY_BEAN_NAME);
+ assertSame(mbf, dispatcherAccessor.getPropertyValue("messageBuilderFactory"));
+ context.close();
}
@Test
@@ -63,6 +69,7 @@ public class PublishSubscribeChannelParserTests {
BroadcastingDispatcher dispatcher = (BroadcastingDispatcher)
accessor.getPropertyValue("dispatcher");
assertTrue((Boolean) new DirectFieldAccessor(dispatcher).getPropertyValue("ignoreFailures"));
+ context.close();
}
@Test
@@ -75,6 +82,7 @@ public class PublishSubscribeChannelParserTests {
BroadcastingDispatcher dispatcher = (BroadcastingDispatcher)
accessor.getPropertyValue("dispatcher");
assertTrue((Boolean) new DirectFieldAccessor(dispatcher).getPropertyValue("applySequence"));
+ context.close();
}
@Test
@@ -93,6 +101,7 @@ public class PublishSubscribeChannelParserTests {
DirectFieldAccessor executorAccessor = new DirectFieldAccessor(executor);
Executor innerExecutor = (Executor) executorAccessor.getPropertyValue("executor");
assertEquals(context.getBean("pool"), innerExecutor);
+ context.close();
}
@Test
@@ -112,6 +121,7 @@ public class PublishSubscribeChannelParserTests {
DirectFieldAccessor executorAccessor = new DirectFieldAccessor(executor);
Executor innerExecutor = (Executor) executorAccessor.getPropertyValue("executor");
assertEquals(context.getBean("pool"), innerExecutor);
+ context.close();
}
@Test
@@ -131,6 +141,7 @@ public class PublishSubscribeChannelParserTests {
DirectFieldAccessor executorAccessor = new DirectFieldAccessor(executor);
Executor innerExecutor = (Executor) executorAccessor.getPropertyValue("executor");
assertEquals(context.getBean("pool"), innerExecutor);
+ context.close();
}
@Test
@@ -143,6 +154,7 @@ public class PublishSubscribeChannelParserTests {
ErrorHandler errorHandler = (ErrorHandler) accessor.getPropertyValue("errorHandler");
assertNotNull(errorHandler);
assertEquals(context.getBean("testErrorHandler"), errorHandler);
+ context.close();
}
}
diff --git a/spring-integration-core/src/test/java/org/springframework/integration/gateway/GatewayInterfaceTests.java b/spring-integration-core/src/test/java/org/springframework/integration/gateway/GatewayInterfaceTests.java
index ff5932a187..1478c53414 100644
--- a/spring-integration-core/src/test/java/org/springframework/integration/gateway/GatewayInterfaceTests.java
+++ b/spring-integration-core/src/test/java/org/springframework/integration/gateway/GatewayInterfaceTests.java
@@ -1,5 +1,5 @@
/*
- * Copyright 2002-2013 the original author or authors.
+ * Copyright 2002-2014 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.
@@ -21,6 +21,7 @@ import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.assertNull;
+import static org.junit.Assert.assertSame;
import static org.junit.Assert.assertThat;
import static org.junit.Assert.assertTrue;
import static org.mockito.Mockito.mock;
@@ -28,6 +29,7 @@ import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import java.lang.reflect.Method;
+import java.util.Map;
import java.util.concurrent.atomic.AtomicBoolean;
import org.hamcrest.Matchers;
@@ -35,12 +37,14 @@ import org.junit.Test;
import org.mockito.Mockito;
import org.springframework.beans.factory.support.DefaultListableBeanFactory;
-import org.springframework.context.ApplicationContext;
+import org.springframework.context.ConfigurableApplicationContext;
import org.springframework.context.support.ClassPathXmlApplicationContext;
import org.springframework.integration.annotation.Gateway;
import org.springframework.integration.annotation.Header;
import org.springframework.integration.channel.DirectChannel;
+import org.springframework.integration.context.IntegrationContextUtils;
import org.springframework.integration.support.MessageBuilder;
+import org.springframework.integration.test.util.TestUtils;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageHandler;
import org.springframework.messaging.MessagingException;
@@ -55,7 +59,7 @@ public class GatewayInterfaceTests {
@Test
public void testWithServiceSuperclassAnnotatedMethod() throws Exception {
- ApplicationContext ac = new ClassPathXmlApplicationContext("GatewayInterfaceTests-context.xml", this.getClass());
+ ConfigurableApplicationContext ac = new ClassPathXmlApplicationContext("GatewayInterfaceTests-context.xml", this.getClass());
DirectChannel channel = ac.getBean("requestChannelFoo", DirectChannel.class);
final Method fooMethod = Foo.class.getMethod("foo", String.class);
final AtomicBoolean called = new AtomicBoolean();
@@ -76,11 +80,16 @@ public class GatewayInterfaceTests {
Bar bar = ac.getBean(Bar.class);
bar.foo("hello");
assertTrue(called.get());
+ Map,?> gateways = TestUtils.getPropertyValue(ac.getBean("&sampleGateway"), "gatewayMap", Map.class);
+ Object mbf = ac.getBean(IntegrationContextUtils.INTEGRATION_MESSAGE_BUILDER_FACTORY_BEAN_NAME);
+ assertSame(mbf, TestUtils.getPropertyValue(gateways.values().iterator().next(),
+ "messageConverter.messageBuilderFactory"));
+ ac.close();
}
@Test
public void testWithServiceSuperclassAnnotatedMethodOverridePE() throws Exception {
- ApplicationContext ac = new ClassPathXmlApplicationContext("GatewayInterfaceTests2-context.xml", this.getClass());
+ ConfigurableApplicationContext ac = new ClassPathXmlApplicationContext("GatewayInterfaceTests2-context.xml", this.getClass());
DirectChannel channel = ac.getBean("requestChannelFoo", DirectChannel.class);
final Method fooMethod = Foo.class.getMethod("foo", String.class);
final AtomicBoolean called = new AtomicBoolean();
@@ -101,22 +110,24 @@ public class GatewayInterfaceTests {
Bar bar = ac.getBean(Bar.class);
bar.foo("hello");
assertTrue(called.get());
+ ac.close();
}
@Test
public void testWithServiceAnnotatedMethod() {
- ApplicationContext ac = new ClassPathXmlApplicationContext("GatewayInterfaceTests-context.xml", this.getClass());
+ ConfigurableApplicationContext ac = new ClassPathXmlApplicationContext("GatewayInterfaceTests-context.xml", this.getClass());
DirectChannel channel = ac.getBean("requestChannelBar", DirectChannel.class);
MessageHandler handler = mock(MessageHandler.class);
channel.subscribe(handler);
Bar bar = ac.getBean(Bar.class);
bar.bar("hello");
verify(handler, times(1)).handleMessage(Mockito.any(Message.class));
+ ac.close();
}
@Test
public void testWithServiceSuperclassUnAnnotatedMethod() throws Exception {
- ApplicationContext ac = new ClassPathXmlApplicationContext("GatewayInterfaceTests-context.xml", this.getClass());
+ ConfigurableApplicationContext ac = new ClassPathXmlApplicationContext("GatewayInterfaceTests-context.xml", this.getClass());
DirectChannel channel = ac.getBean("requestChannelBaz", DirectChannel.class);
final Method bazMethod = Foo.class.getMethod("baz", String.class);
final AtomicBoolean called = new AtomicBoolean();
@@ -137,11 +148,12 @@ public class GatewayInterfaceTests {
Bar bar = ac.getBean(Bar.class);
bar.baz("hello");
assertTrue(called.get());
+ ac.close();
}
@Test
public void testWithServiceUnAnnotatedMethodGlobalHeaderDoesntOverride() throws Exception {
- ApplicationContext ac = new ClassPathXmlApplicationContext("GatewayInterfaceTests-context.xml", this.getClass());
+ ConfigurableApplicationContext ac = new ClassPathXmlApplicationContext("GatewayInterfaceTests-context.xml", this.getClass());
DirectChannel channel = ac.getBean("requestChannelBaz", DirectChannel.class);
final Method quxMethod = Bar.class.getMethod("qux", String.class, String.class);
final AtomicBoolean called = new AtomicBoolean();
@@ -162,55 +174,60 @@ public class GatewayInterfaceTests {
Bar bar = ac.getBean(Bar.class);
bar.qux("hello", "arg1");
assertTrue(called.get());
+ ac.close();
}
@Test
public void testWithServiceCastAsSuperclassAnnotatedMethod() {
- ApplicationContext ac = new ClassPathXmlApplicationContext("GatewayInterfaceTests-context.xml", this.getClass());
+ ConfigurableApplicationContext ac = new ClassPathXmlApplicationContext("GatewayInterfaceTests-context.xml", this.getClass());
DirectChannel channel = ac.getBean("requestChannelFoo", DirectChannel.class);
MessageHandler handler = mock(MessageHandler.class);
channel.subscribe(handler);
Foo foo = ac.getBean(Foo.class);
foo.foo("hello");
verify(handler, times(1)).handleMessage(Mockito.any(Message.class));
+ ac.close();
}
@Test
public void testWithServiceCastAsSuperclassUnAnnotatedMethod() {
- ApplicationContext ac = new ClassPathXmlApplicationContext("GatewayInterfaceTests-context.xml", this.getClass());
+ ConfigurableApplicationContext ac = new ClassPathXmlApplicationContext("GatewayInterfaceTests-context.xml", this.getClass());
DirectChannel channel = ac.getBean("requestChannelBaz", DirectChannel.class);
MessageHandler handler = mock(MessageHandler.class);
channel.subscribe(handler);
Foo foo = ac.getBean(Foo.class);
foo.baz("hello");
verify(handler, times(1)).handleMessage(Mockito.any(Message.class));
+ ac.close();
}
@Test
public void testWithServiceHashcode() throws Exception {
- ApplicationContext ac = new ClassPathXmlApplicationContext("GatewayInterfaceTests-context.xml", this.getClass());
+ ConfigurableApplicationContext ac = new ClassPathXmlApplicationContext("GatewayInterfaceTests-context.xml", this.getClass());
DirectChannel channel = ac.getBean("requestChannelBaz", DirectChannel.class);
MessageHandler handler = mock(MessageHandler.class);
channel.subscribe(handler);
Bar bar = ac.getBean(Bar.class);
assertEquals(bar.hashCode(), ac.getBean(Bar.class).hashCode());
verify(handler, times(0)).handleMessage(Mockito.any(Message.class));
+ ac.close();
}
@Test
public void testWithServiceToString() {
- ApplicationContext ac = new ClassPathXmlApplicationContext("GatewayInterfaceTests-context.xml", this.getClass());
+ ConfigurableApplicationContext ac = new ClassPathXmlApplicationContext("GatewayInterfaceTests-context.xml", this.getClass());
DirectChannel channel = ac.getBean("requestChannelBaz", DirectChannel.class);
MessageHandler handler = mock(MessageHandler.class);
channel.subscribe(handler);
Bar bar = ac.getBean(Bar.class);
bar.toString();
verify(handler, times(0)).handleMessage(Mockito.any(Message.class));
+ ac.close();
}
@Test
public void testWithServiceEquals() throws Exception {
- ApplicationContext ac = new ClassPathXmlApplicationContext("GatewayInterfaceTests-context.xml", this.getClass());
+ ConfigurableApplicationContext ac = new ClassPathXmlApplicationContext("GatewayInterfaceTests-context.xml", this.getClass());
DirectChannel channel = ac.getBean("requestChannelBaz", DirectChannel.class);
MessageHandler handler = mock(MessageHandler.class);
channel.subscribe(handler);
@@ -225,17 +242,19 @@ public class GatewayInterfaceTests {
fb.afterPropertiesSet();
assertFalse(bar.equals(fb.getObject()));
verify(handler, times(0)).handleMessage(Mockito.any(Message.class));
+ ac.close();
}
@Test
public void testWithServiceGetClass() {
- ApplicationContext ac = new ClassPathXmlApplicationContext("GatewayInterfaceTests-context.xml", this.getClass());
+ ConfigurableApplicationContext ac = new ClassPathXmlApplicationContext("GatewayInterfaceTests-context.xml", this.getClass());
DirectChannel channel = ac.getBean("requestChannelBaz", DirectChannel.class);
MessageHandler handler = mock(MessageHandler.class);
channel.subscribe(handler);
Bar bar = ac.getBean(Bar.class);
bar.getClass();
verify(handler, times(0)).handleMessage(Mockito.any(Message.class));
+ ac.close();
}
@Test(expected=IllegalArgumentException.class)
@@ -245,7 +264,7 @@ public class GatewayInterfaceTests {
@Test
public void testWithCustomMapper() {
- ApplicationContext ac = new ClassPathXmlApplicationContext("GatewayInterfaceTests-context.xml", this.getClass());
+ ConfigurableApplicationContext ac = new ClassPathXmlApplicationContext("GatewayInterfaceTests-context.xml", this.getClass());
DirectChannel channel = ac.getBean("requestChannelBaz", DirectChannel.class);
final AtomicBoolean called = new AtomicBoolean();
MessageHandler handler = new MessageHandler() {
@@ -260,11 +279,12 @@ public class GatewayInterfaceTests {
Baz baz = ac.getBean(Baz.class);
baz.baz("hello");
assertTrue(called.get());
+ ac.close();
}
@Test
public void testLateReply() throws Exception {
- ApplicationContext ac = new ClassPathXmlApplicationContext("GatewayInterfaceTests-context.xml", this.getClass());
+ ConfigurableApplicationContext ac = new ClassPathXmlApplicationContext("GatewayInterfaceTests-context.xml", this.getClass());
Bar baz = ac.getBean(Bar.class);
String reply = baz.lateReply("hello");
assertNull(reply);
@@ -273,6 +293,7 @@ public class GatewayInterfaceTests {
assertNotNull(receive);
MessagingException messagingException = (MessagingException) receive.getPayload();
assertThat(messagingException.getMessage(), Matchers.startsWith("Reply message received but the receiving thread has exited due to a timeout"));
+ ac.close();
}
diff --git a/spring-integration-file/src/main/java/org/springframework/integration/file/config/FileToByteArrayTransformerParser.java b/spring-integration-file/src/main/java/org/springframework/integration/file/config/FileToByteArrayTransformerParser.java
index 4821e3f170..696d8d6bdb 100644
--- a/spring-integration-file/src/main/java/org/springframework/integration/file/config/FileToByteArrayTransformerParser.java
+++ b/spring-integration-file/src/main/java/org/springframework/integration/file/config/FileToByteArrayTransformerParser.java
@@ -1,5 +1,5 @@
/*
- * Copyright 2002-2010 the original author or authors.
+ * Copyright 2002-2014 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.
@@ -16,16 +16,19 @@
package org.springframework.integration.file.config;
+import org.springframework.integration.file.transformer.FileToByteArrayTransformer;
+
/**
* Parser for the <file-to-bytes-transformer> element.
- *
+ *
* @author Mark Fisher
+ * @author Gary Russell
*/
public class FileToByteArrayTransformerParser extends AbstractFilePayloadTransformerParser {
@Override
protected String getTransformerClassName() {
- return "org.springframework.integration.file.transformer.FileToByteArrayTransformer";
+ return FileToByteArrayTransformer.class.getName();
}
}
diff --git a/spring-integration-file/src/main/java/org/springframework/integration/file/config/FileToStringTransformerParser.java b/spring-integration-file/src/main/java/org/springframework/integration/file/config/FileToStringTransformerParser.java
index 6e36b1d9d3..86b7c00cb1 100644
--- a/spring-integration-file/src/main/java/org/springframework/integration/file/config/FileToStringTransformerParser.java
+++ b/spring-integration-file/src/main/java/org/springframework/integration/file/config/FileToStringTransformerParser.java
@@ -1,5 +1,5 @@
/*
- * Copyright 2002-2010 the original author or authors.
+ * Copyright 2002-2014 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.
@@ -21,17 +21,19 @@ import org.w3c.dom.Element;
import org.springframework.beans.factory.support.BeanDefinitionBuilder;
import org.springframework.beans.factory.xml.ParserContext;
import org.springframework.integration.config.xml.IntegrationNamespaceUtils;
+import org.springframework.integration.file.transformer.FileToStringTransformer;
/**
* Parser for the <file-to-string-transformer> element.
- *
+ *
* @author Mark Fisher
+ * @author Gary Russell
*/
public class FileToStringTransformerParser extends AbstractFilePayloadTransformerParser {
@Override
protected String getTransformerClassName() {
- return "org.springframework.integration.file.transformer.FileToStringTransformer";
+ return FileToStringTransformer.class.getName();
}
@Override
diff --git a/spring-integration-file/src/main/java/org/springframework/integration/file/transformer/AbstractFilePayloadTransformer.java b/spring-integration-file/src/main/java/org/springframework/integration/file/transformer/AbstractFilePayloadTransformer.java
index d79584d0f1..42629dc3b1 100644
--- a/spring-integration-file/src/main/java/org/springframework/integration/file/transformer/AbstractFilePayloadTransformer.java
+++ b/spring-integration-file/src/main/java/org/springframework/integration/file/transformer/AbstractFilePayloadTransformer.java
@@ -21,6 +21,10 @@ import java.io.File;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
+import org.springframework.beans.BeansException;
+import org.springframework.beans.factory.BeanFactory;
+import org.springframework.beans.factory.BeanFactoryAware;
+import org.springframework.integration.context.IntegrationContextUtils;
import org.springframework.integration.file.FileHeaders;
import org.springframework.integration.support.DefaultMessageBuilderFactory;
import org.springframework.integration.support.MessageBuilderFactory;
@@ -34,7 +38,7 @@ import org.springframework.util.Assert;
*
* @author Mark Fisher
*/
-public abstract class AbstractFilePayloadTransformer implements Transformer {
+public abstract class AbstractFilePayloadTransformer implements Transformer, BeanFactoryAware {
private final Log logger = LogFactory.getLog(this.getClass());
@@ -52,9 +56,9 @@ public abstract class AbstractFilePayloadTransformer implements Transformer {
this.deleteFiles = deleteFiles;
}
- public void setMessageBuilderFactory(MessageBuilderFactory messageBuilderFactory) {
- Assert.notNull(messageBuilderFactory, "'messageBuilderFactory' cannot be null");
- this.messageBuilderFactory = messageBuilderFactory;
+ @Override
+ public void setBeanFactory(BeanFactory beanFactory) throws BeansException {
+ this.messageBuilderFactory = IntegrationContextUtils.getMessageBuilderFactory(beanFactory);
}
@Override
diff --git a/spring-integration-file/src/test/java/org/springframework/integration/file/config/FileToStringTransformerParserTests.java b/spring-integration-file/src/test/java/org/springframework/integration/file/config/FileToStringTransformerParserTests.java
index fc32c1c94d..a6a51ff022 100644
--- a/spring-integration-file/src/test/java/org/springframework/integration/file/config/FileToStringTransformerParserTests.java
+++ b/spring-integration-file/src/test/java/org/springframework/integration/file/config/FileToStringTransformerParserTests.java
@@ -1,5 +1,5 @@
/*
- * Copyright 2002-2009 the original author or authors.
+ * Copyright 2002-2014 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.
@@ -16,21 +16,25 @@
package org.springframework.integration.file.config;
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertSame;
+
import org.junit.Test;
import org.junit.runner.RunWith;
+
import org.springframework.beans.DirectFieldAccessor;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.integration.endpoint.EventDrivenConsumer;
import org.springframework.integration.file.transformer.FileToStringTransformer;
+import org.springframework.integration.support.MessageBuilderFactory;
import org.springframework.integration.transformer.MessageTransformingHandler;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
-import static org.junit.Assert.assertEquals;
-
/**
* @author Mark Fisher
+ * @author Gary Russell
*/
@ContextConfiguration
@RunWith(SpringJUnit4ClassRunner.class)
@@ -40,6 +44,9 @@ public class FileToStringTransformerParserTests {
@Qualifier("transformer")
EventDrivenConsumer endpoint;
+ @Autowired
+ MessageBuilderFactory messageBuilderFactory;
+
@Test
public void checkDeleteFilesValue() {
DirectFieldAccessor endpointAccessor = new DirectFieldAccessor(endpoint);
@@ -50,6 +57,7 @@ public class FileToStringTransformerParserTests {
handlerAccessor.getPropertyValue("transformer");
DirectFieldAccessor transformerAccessor = new DirectFieldAccessor(transformer);
assertEquals(Boolean.TRUE, transformerAccessor.getPropertyValue("deleteFiles"));
+ assertSame(this.messageBuilderFactory, transformerAccessor.getPropertyValue("messageBuilderFactory"));
}
}
diff --git a/spring-integration-ip/src/main/java/org/springframework/integration/ip/config/TcpConnectionFactoryFactoryBean.java b/spring-integration-ip/src/main/java/org/springframework/integration/ip/config/TcpConnectionFactoryFactoryBean.java
index 7c3ac1e451..4901863b54 100644
--- a/spring-integration-ip/src/main/java/org/springframework/integration/ip/config/TcpConnectionFactoryFactoryBean.java
+++ b/spring-integration-ip/src/main/java/org/springframework/integration/ip/config/TcpConnectionFactoryFactoryBean.java
@@ -27,7 +27,6 @@ import org.springframework.context.ApplicationEventPublisherAware;
import org.springframework.context.SmartLifecycle;
import org.springframework.core.serializer.Deserializer;
import org.springframework.core.serializer.Serializer;
-import org.springframework.integration.context.IntegrationContextUtils;
import org.springframework.integration.ip.tcp.connection.AbstractConnectionFactory;
import org.springframework.integration.ip.tcp.connection.AbstractServerConnectionFactory;
import org.springframework.integration.ip.tcp.connection.DefaultTcpNetSSLSocketFactorySupport;
@@ -134,7 +133,7 @@ public class TcpConnectionFactoryFactoryBean extends AbstractFactoryBean,
- OutboundMessageMapper {
+ OutboundMessageMapper,
+ BeanFactoryAware {
protected final Log logger = LogFactory.getLog(this.getClass());
@@ -81,8 +86,9 @@ public class TcpMessageMapper implements
this.applySequence = applySequence;
}
- public void setMessageBuilderFactory(MessageBuilderFactory messageBuilderFactory) {
- this.messageBuilderFactory = messageBuilderFactory;
+ @Override
+ public void setBeanFactory(BeanFactory beanFactory) throws BeansException {
+ this.messageBuilderFactory = IntegrationContextUtils.getMessageBuilderFactory(beanFactory);
}
protected MessageBuilderFactory getMessageBuilderFactory() {
diff --git a/spring-integration-ip/src/main/java/org/springframework/integration/ip/udp/DatagramPacketMessageMapper.java b/spring-integration-ip/src/main/java/org/springframework/integration/ip/udp/DatagramPacketMessageMapper.java
index 47ef078bfb..0796990cda 100644
--- a/spring-integration-ip/src/main/java/org/springframework/integration/ip/udp/DatagramPacketMessageMapper.java
+++ b/spring-integration-ip/src/main/java/org/springframework/integration/ip/udp/DatagramPacketMessageMapper.java
@@ -23,6 +23,10 @@ import java.util.UUID;
import java.util.regex.Matcher;
import java.util.regex.Pattern;
+import org.springframework.beans.BeansException;
+import org.springframework.beans.factory.BeanFactory;
+import org.springframework.beans.factory.BeanFactoryAware;
+import org.springframework.integration.context.IntegrationContextUtils;
import org.springframework.integration.ip.IpHeaders;
import org.springframework.integration.ip.util.RegexUtils;
import org.springframework.integration.mapping.InboundMessageMapper;
@@ -57,7 +61,8 @@ import org.springframework.util.Assert;
* @author Dave Syer
* @since 2.0
*/
-public class DatagramPacketMessageMapper implements InboundMessageMapper, OutboundMessageMapper {
+public class DatagramPacketMessageMapper implements InboundMessageMapper, OutboundMessageMapper,
+ BeanFactoryAware {
private volatile String charset = "UTF-8";
@@ -77,10 +82,6 @@ public class DatagramPacketMessageMapper implements InboundMessageMapper= 0) {
socket.setSoTimeout(this.getSoTimeout());
@@ -365,56 +424,4 @@ public class UnicastSendingMessageHandler extends
}
}
- /**
- * @see java.net.Socket#setReceiveBufferSize(int)
- * @see DatagramSocket#setReceiveBufferSize(int)
- */
- @Override
- public void setSoReceiveBufferSize(int size) {
- this.soReceiveBufferSize = size;
- }
-
- @Override
- public void setLocalAddress(String localAddress) {
- this.localAddress = localAddress;
- }
-
- public void setTaskExecutor(Executor taskExecutor) {
- Assert.notNull(taskExecutor, "'taskExecutor' cannot be null");
- this.taskExecutor = taskExecutor;
- this.taskExecutorSet = true;
- }
-
- /**
- * @param ackCounter the ackCounter to set
- */
- public void setAckCounter(int ackCounter) {
- this.ackCounter = ackCounter;
- }
-
- @Override
- public String getComponentType(){
- return "ip:udp-outbound-channel-adapter";
- }
-
- /**
- * @return the acknowledge
- */
- public boolean isAcknowledge() {
- return acknowledge;
- }
-
- /**
- * @return the ackPort
- */
- public int getAckPort() {
- return ackPort;
- }
-
- /**
- * @return the soReceiveBufferSize
- */
- public int getSoReceiveBufferSize() {
- return soReceiveBufferSize;
- }
}
diff --git a/spring-integration-ip/src/test/java/org/springframework/integration/ip/config/ParserUnitTests.java b/spring-integration-ip/src/test/java/org/springframework/integration/ip/config/ParserUnitTests.java
index b5d1c1de49..c192623d19 100644
--- a/spring-integration-ip/src/test/java/org/springframework/integration/ip/config/ParserUnitTests.java
+++ b/spring-integration-ip/src/test/java/org/springframework/integration/ip/config/ParserUnitTests.java
@@ -1,5 +1,5 @@
/*
- * Copyright 2002-2013 the original author or authors.
+ * Copyright 2002-2014 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.
@@ -42,6 +42,7 @@ import org.springframework.core.serializer.Serializer;
import org.springframework.core.task.TaskExecutor;
import org.springframework.integration.channel.DirectChannel;
import org.springframework.integration.channel.QueueChannel;
+import org.springframework.integration.core.MessagingTemplate;
import org.springframework.integration.endpoint.AbstractEndpoint;
import org.springframework.integration.endpoint.EventDrivenConsumer;
import org.springframework.integration.handler.advice.AbstractRequestHandlerAdvice;
@@ -70,12 +71,12 @@ import org.springframework.integration.ip.udp.MulticastReceivingChannelAdapter;
import org.springframework.integration.ip.udp.MulticastSendingMessageHandler;
import org.springframework.integration.ip.udp.UnicastReceivingChannelAdapter;
import org.springframework.integration.ip.udp.UnicastSendingMessageHandler;
-import org.springframework.messaging.support.GenericMessage;
+import org.springframework.integration.support.MessageBuilderFactory;
import org.springframework.integration.test.util.TestUtils;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageChannel;
import org.springframework.messaging.MessageHandler;
-import org.springframework.integration.core.MessagingTemplate;
+import org.springframework.messaging.support.GenericMessage;
import org.springframework.scheduling.TaskScheduler;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
@@ -274,6 +275,9 @@ public class ParserUnitTests {
@Autowired
QueueChannel eventChannel;
+ @Autowired
+ MessageBuilderFactory messageBuilderFactory;
+
private static volatile int adviceCalled;
@Test
@@ -295,6 +299,7 @@ public class ParserUnitTests {
assertFalse((Boolean)mapperAccessor.getPropertyValue("lookupHost"));
assertFalse(TestUtils.getPropertyValue(udpIn, "autoStartup", Boolean.class));
assertEquals(1234, dfa.getPropertyValue("phase"));
+ assertSame(this.messageBuilderFactory, mapperAccessor.getPropertyValue("messageBuilderFactory"));
}
@Test
@@ -313,12 +318,14 @@ public class ParserUnitTests {
DatagramPacketMessageMapper mapper = (DatagramPacketMessageMapper) dfa.getPropertyValue("mapper");
DirectFieldAccessor mapperAccessor = new DirectFieldAccessor(mapper);
assertTrue((Boolean)mapperAccessor.getPropertyValue("lookupHost"));
+ assertSame(this.messageBuilderFactory, mapperAccessor.getPropertyValue("messageBuilderFactory"));
}
@Test
public void testInTcp() {
DirectFieldAccessor dfa = new DirectFieldAccessor(tcpIn);
assertSame(cfS1, dfa.getPropertyValue("serverConnectionFactory"));
+ assertSame(this.messageBuilderFactory, TestUtils.getPropertyValue(cfS1, "mapper.messageBuilderFactory"));
assertEquals("testInTcp",tcpIn.getComponentName());
assertEquals("ip:tcp-inbound-channel-adapter", tcpIn.getComponentType());
assertEquals(errorChannel, dfa.getPropertyValue("errorChannel"));
@@ -345,8 +352,8 @@ public class ParserUnitTests {
@Test
public void testInTcpNioSSLDefaultConfig() {
assertFalse(cfS1Nio.isLookupHost());
- assertTrue((Boolean) TestUtils.getPropertyValue(
- TestUtils.getPropertyValue(cfS1Nio, "mapper"), "applySequence"));
+ assertTrue((Boolean) TestUtils.getPropertyValue(cfS1Nio, "mapper.applySequence"));
+ assertSame(this.messageBuilderFactory, TestUtils.getPropertyValue(cfS1Nio, "mapper.messageBuilderFactory"));
Object connectionSupport = TestUtils.getPropertyValue(cfS1Nio, "tcpNioConnectionSupport");
assertTrue(connectionSupport instanceof DefaultTcpNioSSLConnectionSupport);
assertNotNull(TestUtils.getPropertyValue(connectionSupport, "sslContext"));
@@ -374,6 +381,7 @@ public class ParserUnitTests {
assertEquals(23, dfa.getPropertyValue("order"));
assertEquals("testOutUdp",udpOut.getComponentName());
assertEquals("ip:udp-outbound-channel-adapter", udpOut.getComponentType());
+ assertSame(this.messageBuilderFactory, TestUtils.getPropertyValue(mapper, "messageBuilderFactory"));
}
@Test
@@ -395,6 +403,7 @@ public class ParserUnitTests {
assertEquals(54, dfa.getPropertyValue("soTimeout"));
assertEquals(55, dfa.getPropertyValue("timeToLive"));
assertEquals(12, dfa.getPropertyValue("order"));
+ assertSame(this.messageBuilderFactory, TestUtils.getPropertyValue(mapper, "messageBuilderFactory"));
}
@Test
diff --git a/spring-integration-jms/src/main/java/org/springframework/integration/jms/SubscribableJmsChannel.java b/spring-integration-jms/src/main/java/org/springframework/integration/jms/SubscribableJmsChannel.java
index a8bda09d2e..e479fa0180 100644
--- a/spring-integration-jms/src/main/java/org/springframework/integration/jms/SubscribableJmsChannel.java
+++ b/spring-integration-jms/src/main/java/org/springframework/integration/jms/SubscribableJmsChannel.java
@@ -102,7 +102,9 @@ public class SubscribableJmsChannel extends AbstractJmsChannel implements Subscr
private void configureDispatcher(boolean isPubSub) {
if (isPubSub) {
- this.dispatcher = new BroadcastingDispatcher(true);
+ BroadcastingDispatcher broadcastingDispatcher = new BroadcastingDispatcher(true);
+ broadcastingDispatcher.setBeanFactory(this.getBeanFactory());
+ this.dispatcher = broadcastingDispatcher;
}
else {
UnicastingDispatcher unicastingDispatcher = new UnicastingDispatcher();
diff --git a/spring-integration-jms/src/main/java/org/springframework/integration/jms/config/JmsChannelFactoryBean.java b/spring-integration-jms/src/main/java/org/springframework/integration/jms/config/JmsChannelFactoryBean.java
index 503e05387f..5294707a15 100644
--- a/spring-integration-jms/src/main/java/org/springframework/integration/jms/config/JmsChannelFactoryBean.java
+++ b/spring-integration-jms/src/main/java/org/springframework/integration/jms/config/JmsChannelFactoryBean.java
@@ -337,6 +337,9 @@ public class JmsChannelFactoryBean extends AbstractFactoryBean message) throws MessagingException {
assertEquals(m.getPayload(), message.getPayload());
latch.countDown();
@@ -84,7 +89,7 @@ public class RedisChannelParserTests extends RedisAvailableTests{
redisChannel.send(m);
assertTrue(latch.await(5, TimeUnit.SECONDS));
- context.stop();
+ context.close();
}
}
diff --git a/spring-integration-redis/src/test/java/org/springframework/integration/redis/config/RedisInboundChannelAdapterParserTests.java b/spring-integration-redis/src/test/java/org/springframework/integration/redis/config/RedisInboundChannelAdapterParserTests.java
index dc1b29efe1..8e724b8f50 100644
--- a/spring-integration-redis/src/test/java/org/springframework/integration/redis/config/RedisInboundChannelAdapterParserTests.java
+++ b/spring-integration-redis/src/test/java/org/springframework/integration/redis/config/RedisInboundChannelAdapterParserTests.java
@@ -1,5 +1,5 @@
/*
- * Copyright 2002-2013 the original author or authors.
+ * Copyright 2002-2014 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.
@@ -33,6 +33,7 @@ import org.springframework.context.ApplicationContext;
import org.springframework.data.redis.connection.RedisConnectionFactory;
import org.springframework.data.redis.listener.RedisMessageListenerContainer;
import org.springframework.integration.channel.QueueChannel;
+import org.springframework.integration.context.IntegrationContextUtils;
import org.springframework.integration.redis.inbound.RedisInboundChannelAdapter;
import org.springframework.integration.redis.rules.RedisAvailable;
import org.springframework.integration.redis.rules.RedisAvailableTests;
@@ -77,6 +78,8 @@ public class RedisInboundChannelAdapterParserTests extends RedisAvailableTests {
Object bean = context.getBean("withoutSerializer.adapter");
assertNotNull(bean);
assertNull(TestUtils.getPropertyValue(bean, "serializer"));
+ Object mbf = context.getBean(IntegrationContextUtils.INTEGRATION_MESSAGE_BUILDER_FACTORY_BEAN_NAME);
+ assertSame(mbf, TestUtils.getPropertyValue(bean, "messageConverter.messageBuilderFactory"));
}
@Test
diff --git a/spring-integration-redis/src/test/java/org/springframework/integration/redis/config/RedisOutboundChannelAdapterParserTests.java b/spring-integration-redis/src/test/java/org/springframework/integration/redis/config/RedisOutboundChannelAdapterParserTests.java
index 159d617f5f..bbe573faf8 100644
--- a/spring-integration-redis/src/test/java/org/springframework/integration/redis/config/RedisOutboundChannelAdapterParserTests.java
+++ b/spring-integration-redis/src/test/java/org/springframework/integration/redis/config/RedisOutboundChannelAdapterParserTests.java
@@ -18,6 +18,7 @@ package org.springframework.integration.redis.config;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertNotNull;
+import static org.junit.Assert.assertSame;
import org.junit.Test;
import org.junit.runner.RunWith;
@@ -28,6 +29,7 @@ import org.springframework.context.ApplicationContext;
import org.springframework.data.redis.listener.RedisMessageListenerContainer;
import org.springframework.expression.Expression;
import org.springframework.integration.channel.QueueChannel;
+import org.springframework.integration.context.IntegrationContextUtils;
import org.springframework.integration.endpoint.EventDrivenConsumer;
import org.springframework.integration.redis.inbound.RedisInboundChannelAdapter;
import org.springframework.integration.redis.outbound.RedisPublishingMessageHandler;
@@ -73,6 +75,8 @@ public class RedisOutboundChannelAdapterParserTests extends RedisAvailableTests
Object converterBean = context.getBean("testConverter");
assertEquals(converterBean, accessor.getPropertyValue("messageConverter"));
assertEquals(context.getBean("serializer"), accessor.getPropertyValue("serializer"));
+ Object mbf = context.getBean(IntegrationContextUtils.INTEGRATION_MESSAGE_BUILDER_FACTORY_BEAN_NAME);
+ assertSame(mbf, TestUtils.getPropertyValue(handler, "messageConverter.messageBuilderFactory"));
}
@Test
diff --git a/spring-integration-redis/src/test/java/org/springframework/integration/redis/config/RedisStoreInboundChannelAdapterParserTests.java b/spring-integration-redis/src/test/java/org/springframework/integration/redis/config/RedisStoreInboundChannelAdapterParserTests.java
index a2753bbd22..7fd9ab535b 100644
--- a/spring-integration-redis/src/test/java/org/springframework/integration/redis/config/RedisStoreInboundChannelAdapterParserTests.java
+++ b/spring-integration-redis/src/test/java/org/springframework/integration/redis/config/RedisStoreInboundChannelAdapterParserTests.java
@@ -1,5 +1,5 @@
/*
- * Copyright 2002-2012 the original author or authors.
+ * Copyright 2002-2014 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.
@@ -22,6 +22,7 @@ import static org.junit.Assert.assertTrue;
import org.junit.Test;
import org.junit.runner.RunWith;
+
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.parsing.BeanDefinitionParsingException;
import org.springframework.context.ApplicationContext;
@@ -70,7 +71,7 @@ public class RedisStoreInboundChannelAdapterParserTests {
@Test(expected=BeanDefinitionParsingException.class)
public void testTemplateAndCfMutualExclusivity(){
- new ClassPathXmlApplicationContext("inbound-template-cf-fail.xml", this.getClass());
+ new ClassPathXmlApplicationContext("inbound-template-cf-fail.xml", this.getClass()).close();
}
}