From 93f6568c9c29b2a5dbd5caf168a7de8f94aaa1f5 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Fri, 21 Jan 2011 11:06:48 -0500 Subject: [PATCH] INT-1758 refactored AbsractPollingEndpoint to remove deopendency on PollerMetadata, fixed tests --- .../config/ConsumerEndpointFactoryBean.java | 6 ++- ...ourcePollingChannelAdapterFactoryBean.java | 5 +- .../endpoint/AbstractPollingEndpoint.java | 37 ++++++++++---- .../gateway/MessagingGatewaySupport.java | 2 - .../scheduling/PollerMetadata.java | 4 +- .../ApplicationContextMessageBusTests.java | 5 -- .../integration/bus/messageBusTests.xml | 10 ++-- .../InboundChannelAdapterExpressionTests.java | 8 ++-- .../core/MessagingTemplateTests.java | 6 +-- .../dispatcher/PollingTransactionTests.java | 5 +- ...aluatingMessageSourceIntegrationTests.java | 8 ++-- .../PollingConsumerEndpointTests.java | 36 +++++--------- .../endpoint/PollingEndpointStub.java | 5 +- .../endpoint/PollingLifecycleTests.java | 6 +-- .../pollingEndpointErrorHandlingTests.xml | 6 +-- .../MethodInvokingMessageHandlerTests.java | 5 +- .../ByteStreamWritingMessageHandlerTests.java | 48 +++++++------------ ...acterStreamWritingMessageHandlerTests.java | 48 +++++++------------ .../WebServiceOutboundGatewayParserTests.java | 4 +- 19 files changed, 100 insertions(+), 154 deletions(-) diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/ConsumerEndpointFactoryBean.java b/spring-integration-core/src/main/java/org/springframework/integration/config/ConsumerEndpointFactoryBean.java index 87f94f2636..7567c35ca3 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/ConsumerEndpointFactoryBean.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/ConsumerEndpointFactoryBean.java @@ -158,7 +158,11 @@ public class ConsumerEndpointFactoryBean Assert.notNull(this.pollerMetadata, "No poller has been defined for endpoint '" + this.beanName + "', and no default poller is available within the context."); } - pollingConsumer.setPollerMetadata(this.pollerMetadata); + pollingConsumer.setTaskExecutor(this.pollerMetadata.getTaskExecutor()); + pollingConsumer.setTrigger(this.pollerMetadata.getTrigger()); + pollingConsumer.setAdviceChain(this.pollerMetadata.getAdviceChain()); + pollingConsumer.setMaxMessagesPerPoll(this.pollerMetadata.getMaxMessagesPerPoll()); + pollingConsumer.setReceiveTimeout(this.pollerMetadata.getReceiveTimeout()); pollingConsumer.setBeanClassLoader(beanClassLoader); pollingConsumer.setBeanFactory(beanFactory); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/SourcePollingChannelAdapterFactoryBean.java b/spring-integration-core/src/main/java/org/springframework/integration/config/SourcePollingChannelAdapterFactoryBean.java index bd70333b0d..925ecba2cc 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/SourcePollingChannelAdapterFactoryBean.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/SourcePollingChannelAdapterFactoryBean.java @@ -129,7 +129,10 @@ public class SourcePollingChannelAdapterFactoryBean implements FactoryBean adviceChain; private volatile ClassLoader beanClassLoader = ClassUtils.getDefaultClassLoader(); @@ -56,17 +60,30 @@ public abstract class AbstractPollingEndpoint extends AbstractEndpoint implement private volatile Runnable poller; private volatile boolean initialized; - + + private volatile long maxMessagesPerPoll = -1; + private final Object initializationMonitor = new Object(); public AbstractPollingEndpoint() { this.setPhase(Integer.MAX_VALUE); } + + public void setTaskExecutor(Executor taskExecutor) { + this.taskExecutor = (taskExecutor != null ? taskExecutor : new SyncTaskExecutor()); + } + public void setTrigger(Trigger trigger) { + this.trigger = (trigger != null ? trigger : new PeriodicTrigger(10)); + } - public void setPollerMetadata(PollerMetadata pollerMetadata) { - this.pollerMetadata = pollerMetadata; + public void setAdviceChain(List adviceChain) { + this.adviceChain = adviceChain; + } + + public void setMaxMessagesPerPoll(long maxMessagesPerPoll) { + this.maxMessagesPerPoll = maxMessagesPerPoll; } public void setErrorHandler(ErrorHandler errorHandler) { @@ -83,8 +100,8 @@ public abstract class AbstractPollingEndpoint extends AbstractEndpoint implement if (this.initialized) { return; } - Assert.notNull(this.pollerMetadata.getTrigger(), "Trigger is required"); - Executor providedExecutor = this.pollerMetadata.getTaskExecutor(); + Assert.notNull(this.trigger, "Trigger is required"); + Executor providedExecutor = this.taskExecutor; if (providedExecutor != null) { this.taskExecutor = providedExecutor; } @@ -117,7 +134,7 @@ public abstract class AbstractPollingEndpoint extends AbstractEndpoint implement } }; - List adviceChain = this.pollerMetadata.getAdviceChain(); + List adviceChain = this.adviceChain; if (!CollectionUtils.isEmpty(adviceChain)) { ProxyFactory proxyFactory = new ProxyFactory(pollingTask); if (!CollectionUtils.isEmpty(adviceChain)) { @@ -140,7 +157,7 @@ public abstract class AbstractPollingEndpoint extends AbstractEndpoint implement } Assert.state(this.getTaskScheduler() != null, "unable to start polling, no taskScheduler available"); - this.runningTask = this.getTaskScheduler().schedule(this.poller, this.pollerMetadata.getTrigger()); + this.runningTask = this.getTaskScheduler().schedule(this.poller, this.trigger); } @Override // guarded by super#lifecycleLock @@ -161,7 +178,7 @@ public abstract class AbstractPollingEndpoint extends AbstractEndpoint implement */ private class Poller implements Runnable { - private final long maxMessagesPerPoll = pollerMetadata.getMaxMessagesPerPoll(); + //private final long maxMessagesPerPoll = pollerMetadata.getMaxMessagesPerPoll(); private final Callable pollingTask; 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 9820a0d546..07af8f8e5b 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 @@ -32,7 +32,6 @@ import org.springframework.integration.history.TrackableComponent; import org.springframework.integration.mapping.InboundMessageMapper; import org.springframework.integration.mapping.OutboundMessageMapper; import org.springframework.integration.message.ErrorMessage; -import org.springframework.integration.scheduling.PollerMetadata; import org.springframework.integration.support.MessageBuilder; import org.springframework.integration.support.converter.SimpleMessageConverter; import org.springframework.util.Assert; @@ -292,7 +291,6 @@ public abstract class MessagingGatewaySupport extends AbstractEndpoint implement else if (this.replyChannel instanceof PollableChannel) { PollingConsumer endpoint = new PollingConsumer( (PollableChannel) this.replyChannel, handler); - endpoint.setPollerMetadata(new PollerMetadata()); endpoint.setBeanFactory(this.getBeanFactory()); endpoint.setReceiveTimeout(this.replyTimeout); endpoint.afterPropertiesSet(); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/scheduling/PollerMetadata.java b/spring-integration-core/src/main/java/org/springframework/integration/scheduling/PollerMetadata.java index 0aa556ef29..bf85e50d4b 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/scheduling/PollerMetadata.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/scheduling/PollerMetadata.java @@ -20,8 +20,8 @@ import java.util.List; import java.util.concurrent.Executor; import org.aopalliance.aop.Advice; + import org.springframework.scheduling.Trigger; -import org.springframework.scheduling.support.PeriodicTrigger; import org.springframework.util.ErrorHandler; /** @@ -32,7 +32,7 @@ public class PollerMetadata { public static final int MAX_MESSAGES_UNBOUNDED = -1; - private volatile Trigger trigger = new PeriodicTrigger(10); + private volatile Trigger trigger; private volatile long maxMessagesPerPoll = MAX_MESSAGES_UNBOUNDED; diff --git a/spring-integration-core/src/test/java/org/springframework/integration/bus/ApplicationContextMessageBusTests.java b/spring-integration-core/src/test/java/org/springframework/integration/bus/ApplicationContextMessageBusTests.java index 979f56fe28..de35a1b92f 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/bus/ApplicationContextMessageBusTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/bus/ApplicationContextMessageBusTests.java @@ -69,7 +69,6 @@ public class ApplicationContextMessageBusTests { handler.setBeanFactory(context); handler.afterPropertiesSet(); PollingConsumer endpoint = new PollingConsumer(sourceChannel, handler); - endpoint.setPollerMetadata(new PollerMetadata()); endpoint.setBeanFactory(mock(BeanFactory.class)); context.registerEndpoint("testEndpoint", endpoint); context.refresh(); @@ -127,10 +126,8 @@ public class ApplicationContextMessageBusTests { handler1.setOutputChannel(outputChannel1); handler2.setOutputChannel(outputChannel2); PollingConsumer endpoint1 = new PollingConsumer(inputChannel, handler1); - endpoint1.setPollerMetadata(new PollerMetadata()); endpoint1.setBeanFactory(mock(BeanFactory.class)); PollingConsumer endpoint2 = new PollingConsumer(inputChannel, handler2); - endpoint2.setPollerMetadata(new PollerMetadata()); endpoint2.setBeanFactory(mock(BeanFactory.class)); context.registerEndpoint("testEndpoint1", endpoint1); context.registerEndpoint("testEndpoint2", endpoint2); @@ -194,7 +191,6 @@ public class ApplicationContextMessageBusTests { channelAdapter.setSource(new FailingSource(latch)); PollerMetadata pollerMetadata = new PollerMetadata(); pollerMetadata.setTrigger(new PeriodicTrigger(1000)); - channelAdapter.setPollerMetadata(pollerMetadata); channelAdapter.setOutputChannel(outputChannel); context.registerEndpoint("testChannel", channelAdapter); context.refresh(); @@ -222,7 +218,6 @@ public class ApplicationContextMessageBusTests { } }; PollingConsumer endpoint = new PollingConsumer(errorChannel, handler); - endpoint.setPollerMetadata(new PollerMetadata()); endpoint.setBeanFactory(mock(BeanFactory.class)); context.registerEndpoint("testEndpoint", endpoint); context.refresh(); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/bus/messageBusTests.xml b/spring-integration-core/src/test/java/org/springframework/integration/bus/messageBusTests.xml index e0f4d2b28d..89b1ecb18f 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/bus/messageBusTests.xml +++ b/spring-integration-core/src/test/java/org/springframework/integration/bus/messageBusTests.xml @@ -17,13 +17,9 @@ class="org.springframework.integration.endpoint.PollingConsumer"> - - diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/InboundChannelAdapterExpressionTests.java b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/InboundChannelAdapterExpressionTests.java index 6d1e7b0acc..6777d49477 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/InboundChannelAdapterExpressionTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/InboundChannelAdapterExpressionTests.java @@ -53,7 +53,7 @@ public class InboundChannelAdapterExpressionTests { SourcePollingChannelAdapter adapter = context.getBean("fixedDelayProducer", SourcePollingChannelAdapter.class); assertFalse(adapter.isAutoStartup()); DirectFieldAccessor adapterAccessor = new DirectFieldAccessor(adapter); - Trigger trigger = TestUtils.getPropertyValue(adapter, "pollerMetadata.trigger", Trigger.class); + Trigger trigger = TestUtils.getPropertyValue(adapter, "trigger", Trigger.class); assertEquals(PeriodicTrigger.class, trigger.getClass()); DirectFieldAccessor triggerAccessor = new DirectFieldAccessor(trigger); assertEquals(1234L, triggerAccessor.getPropertyValue("period")); @@ -68,7 +68,7 @@ public class InboundChannelAdapterExpressionTests { SourcePollingChannelAdapter adapter = context.getBean("fixedRateProducer", SourcePollingChannelAdapter.class); assertFalse(adapter.isAutoStartup()); DirectFieldAccessor adapterAccessor = new DirectFieldAccessor(adapter); - Trigger trigger = TestUtils.getPropertyValue(adapter, "pollerMetadata.trigger", Trigger.class); + Trigger trigger = TestUtils.getPropertyValue(adapter, "trigger", Trigger.class); assertEquals(PeriodicTrigger.class, trigger.getClass()); DirectFieldAccessor triggerAccessor = new DirectFieldAccessor(trigger); assertEquals(5678L, triggerAccessor.getPropertyValue("period")); @@ -83,7 +83,7 @@ public class InboundChannelAdapterExpressionTests { SourcePollingChannelAdapter adapter = context.getBean("cronProducer", SourcePollingChannelAdapter.class); assertFalse(adapter.isAutoStartup()); DirectFieldAccessor adapterAccessor = new DirectFieldAccessor(adapter); - Trigger trigger = TestUtils.getPropertyValue(adapter, "pollerMetadata.trigger", Trigger.class); + Trigger trigger = TestUtils.getPropertyValue(adapter, "trigger", Trigger.class); assertEquals(CronTrigger.class, trigger.getClass()); assertEquals("7 6 5 4 3 ?", new DirectFieldAccessor(new DirectFieldAccessor( trigger).getPropertyValue("sequenceGenerator")).getPropertyValue("expression")); @@ -97,7 +97,7 @@ public class InboundChannelAdapterExpressionTests { SourcePollingChannelAdapter adapter = context.getBean("triggerRefProducer", SourcePollingChannelAdapter.class); assertTrue(adapter.isAutoStartup()); DirectFieldAccessor adapterAccessor = new DirectFieldAccessor(adapter); - Trigger trigger = TestUtils.getPropertyValue(adapter, "pollerMetadata.trigger", Trigger.class); + Trigger trigger = TestUtils.getPropertyValue(adapter, "trigger", Trigger.class); assertEquals(context.getBean("customTrigger"), trigger); assertEquals(context.getBean("triggerRefChannel"), adapterAccessor.getPropertyValue("outputChannel")); Expression expression = TestUtils.getPropertyValue(adapter, "source.expression", Expression.class); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/core/MessagingTemplateTests.java b/spring-integration-core/src/test/java/org/springframework/integration/core/MessagingTemplateTests.java index d521637b9d..669b81c302 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/core/MessagingTemplateTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/core/MessagingTemplateTests.java @@ -29,6 +29,7 @@ import java.util.concurrent.TimeUnit; import org.junit.After; import org.junit.Before; import org.junit.Test; + import org.springframework.beans.factory.support.DefaultListableBeanFactory; import org.springframework.context.support.StaticApplicationContext; import org.springframework.integration.Message; @@ -41,14 +42,12 @@ import org.springframework.integration.handler.AbstractReplyProducingMessageHand import org.springframework.integration.mapping.InboundMessageMapper; import org.springframework.integration.mapping.OutboundMessageMapper; import org.springframework.integration.message.GenericMessage; -import org.springframework.integration.scheduling.PollerMetadata; import org.springframework.integration.support.MessageBuilder; import org.springframework.integration.support.channel.ChannelResolutionException; import org.springframework.integration.support.channel.ChannelResolver; import org.springframework.integration.support.converter.SimpleMessageConverter; import org.springframework.integration.test.util.TestUtils; import org.springframework.integration.test.util.TestUtils.TestApplicationContext; -import org.springframework.scheduling.support.PeriodicTrigger; /** * @author Mark Fisher @@ -65,9 +64,6 @@ public class MessagingTemplateTests { this.requestChannel = new QueueChannel(); context.registerChannel("requestChannel", requestChannel); PollingConsumer endpoint = new PollingConsumer(requestChannel, new TestHandler()); - PollerMetadata pollerMetadata = new PollerMetadata(); - pollerMetadata.setTrigger(new PeriodicTrigger(10)); - endpoint.setPollerMetadata(pollerMetadata); context.registerEndpoint("testEndpoint", endpoint); context.refresh(); } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/dispatcher/PollingTransactionTests.java b/spring-integration-core/src/test/java/org/springframework/integration/dispatcher/PollingTransactionTests.java index 549ee8a439..996af8339d 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/dispatcher/PollingTransactionTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/dispatcher/PollingTransactionTests.java @@ -28,6 +28,7 @@ import org.aopalliance.aop.Advice; import org.aopalliance.intercept.MethodInterceptor; import org.aopalliance.intercept.MethodInvocation; import org.junit.Test; + import org.springframework.aop.Advisor; import org.springframework.aop.framework.Advised; import org.springframework.aop.support.DefaultPointcutAdvisor; @@ -37,7 +38,6 @@ import org.springframework.integration.MessageChannel; import org.springframework.integration.core.PollableChannel; import org.springframework.integration.endpoint.PollingConsumer; import org.springframework.integration.message.GenericMessage; -import org.springframework.integration.scheduling.PollerMetadata; import org.springframework.integration.test.util.TestUtils; import org.springframework.integration.util.TestTransactionManager; import org.springframework.transaction.IllegalTransactionStateException; @@ -75,8 +75,7 @@ public class PollingTransactionTests { "transactionTests.xml", this.getClass()); PollingConsumer advicedPoller = context.getBean("advicedSa", PollingConsumer.class); - PollerMetadata pollerMetedata = TestUtils.getPropertyValue(advicedPoller, "pollerMetadata",PollerMetadata.class); - List adviceChain = TestUtils.getPropertyValue(pollerMetedata, "adviceChain",List.class); + List adviceChain = TestUtils.getPropertyValue(advicedPoller, "adviceChain",List.class); assertEquals(3, adviceChain.size()); Runnable poller = TestUtils.getPropertyValue(advicedPoller, "poller", Runnable.class); Callable pollingTask = TestUtils.getPropertyValue(poller, "pollingTask", Callable.class); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/ExpressionEvaluatingMessageSourceIntegrationTests.java b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/ExpressionEvaluatingMessageSourceIntegrationTests.java index 23a688d5fe..88817f08b0 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/ExpressionEvaluatingMessageSourceIntegrationTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/ExpressionEvaluatingMessageSourceIntegrationTests.java @@ -25,13 +25,13 @@ import java.util.Map; import java.util.concurrent.atomic.AtomicInteger; import org.junit.Test; + import org.springframework.expression.Expression; import org.springframework.expression.common.LiteralExpression; import org.springframework.expression.spel.standard.SpelExpressionParser; import org.springframework.integration.Message; import org.springframework.integration.channel.QueueChannel; import org.springframework.integration.config.ExpressionFactoryBean; -import org.springframework.integration.scheduling.PollerMetadata; import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler; import org.springframework.scheduling.support.PeriodicTrigger; import org.springframework.util.ErrorHandler; @@ -62,10 +62,8 @@ public class ExpressionEvaluatingMessageSourceIntegrationTests { SourcePollingChannelAdapter adapter = new SourcePollingChannelAdapter(); adapter.setSource(source); adapter.setTaskScheduler(scheduler); - PollerMetadata pollerMetadata = new PollerMetadata(); - pollerMetadata.setMaxMessagesPerPoll(3); - pollerMetadata.setTrigger(new PeriodicTrigger(60000)); - adapter.setPollerMetadata(pollerMetadata); + adapter.setMaxMessagesPerPoll(3); + adapter.setTrigger(new PeriodicTrigger(60000)); adapter.setOutputChannel(channel); adapter.setErrorHandler(new ErrorHandler() { public void handleError(Throwable t) { diff --git a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollingConsumerEndpointTests.java b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollingConsumerEndpointTests.java index d1f0ccea6c..2effb58eaa 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollingConsumerEndpointTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollingConsumerEndpointTests.java @@ -41,7 +41,6 @@ import org.springframework.integration.MessageRejectedException; import org.springframework.integration.core.MessageHandler; import org.springframework.integration.core.PollableChannel; import org.springframework.integration.message.GenericMessage; -import org.springframework.integration.scheduling.PollerMetadata; import org.springframework.scheduling.Trigger; import org.springframework.scheduling.TriggerContext; import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler; @@ -81,9 +80,7 @@ public class PollingConsumerEndpointTests { taskScheduler.setPoolSize(5); endpoint.setErrorHandler(errorHandler); endpoint.setTaskScheduler(taskScheduler); - PollerMetadata pollerMetadata = new PollerMetadata(); - pollerMetadata.setTrigger(trigger); - endpoint.setPollerMetadata(pollerMetadata); + endpoint.setTrigger(trigger); endpoint.setBeanFactory(mock(BeanFactory.class)); endpoint.setReceiveTimeout(-1); endpoint.afterPropertiesSet(); @@ -102,10 +99,8 @@ public class PollingConsumerEndpointTests { expect(channelMock.receive()).andReturn(message); expectLastCall(); replay(channelMock); - PollerMetadata pollerMetadata = new PollerMetadata(); - pollerMetadata.setMaxMessagesPerPoll(1); - pollerMetadata.setTrigger(trigger); - endpoint.setPollerMetadata(pollerMetadata); + endpoint.setMaxMessagesPerPoll(1); + endpoint.setTrigger(trigger); endpoint.start(); trigger.await(); endpoint.stop(); @@ -117,10 +112,8 @@ public class PollingConsumerEndpointTests { public void multipleMessages() { expect(channelMock.receive()).andReturn(message).times(5); replay(channelMock); - PollerMetadata pollerMetadata = new PollerMetadata(); - pollerMetadata.setMaxMessagesPerPoll(5); - pollerMetadata.setTrigger(trigger); - endpoint.setPollerMetadata(pollerMetadata); + endpoint.setMaxMessagesPerPoll(5); + endpoint.setTrigger(trigger); endpoint.start(); trigger.await(); endpoint.stop(); @@ -133,10 +126,8 @@ public class PollingConsumerEndpointTests { expect(channelMock.receive()).andReturn(message).times(5); expect(channelMock.receive()).andReturn(null); replay(channelMock); - PollerMetadata pollerMetadata = new PollerMetadata(); - pollerMetadata.setMaxMessagesPerPoll(6); - pollerMetadata.setTrigger(trigger); - endpoint.setPollerMetadata(pollerMetadata); + endpoint.setMaxMessagesPerPoll(6); + endpoint.setTrigger(trigger); endpoint.start(); trigger.await(); endpoint.stop(); @@ -169,11 +160,8 @@ public class PollingConsumerEndpointTests { public void droppedMessage_onePerPoll() throws Throwable { expect(channelMock.receive()).andReturn(badMessage).times(1); replay(channelMock); - PollerMetadata pollerMetadata = new PollerMetadata(); - pollerMetadata.setMaxMessagesPerPoll(10); - pollerMetadata.setTrigger(trigger); - endpoint.setPollerMetadata(pollerMetadata); - //endpoint.setErrorHandler(null); + endpoint.setMaxMessagesPerPoll(10); + endpoint.setTrigger(trigger); endpoint.start(); trigger.await(); endpoint.stop(); @@ -201,10 +189,8 @@ public class PollingConsumerEndpointTests { expectLastCall(); replay(channelMock); endpoint.setReceiveTimeout(1); - PollerMetadata pollerMetadata = new PollerMetadata(); - pollerMetadata.setMaxMessagesPerPoll(1); - pollerMetadata.setTrigger(trigger); - endpoint.setPollerMetadata(pollerMetadata); + endpoint.setMaxMessagesPerPoll(1); + endpoint.setTrigger(trigger); endpoint.start(); trigger.await(); endpoint.stop(); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollingEndpointStub.java b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollingEndpointStub.java index df0f625f57..831c9b9a5c 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollingEndpointStub.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollingEndpointStub.java @@ -16,7 +16,6 @@ package org.springframework.integration.endpoint; -import org.springframework.integration.scheduling.PollerMetadata; import org.springframework.scheduling.support.PeriodicTrigger; /** @@ -25,9 +24,7 @@ import org.springframework.scheduling.support.PeriodicTrigger; public class PollingEndpointStub extends AbstractPollingEndpoint { public PollingEndpointStub() { - PollerMetadata pollerMetadata = new PollerMetadata(); - pollerMetadata.setTrigger(new PeriodicTrigger(500)); - this.setPollerMetadata(pollerMetadata); + this.setTrigger(new PeriodicTrigger(500)); } @Override diff --git a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollingLifecycleTests.java b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollingLifecycleTests.java index 0607413145..b770bd931e 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollingLifecycleTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollingLifecycleTests.java @@ -68,9 +68,7 @@ public class PollingLifecycleTests { } }); PollingConsumer consumer = new PollingConsumer(channel, handler); - PollerMetadata pollerMetadata = new PollerMetadata(); - pollerMetadata.setTrigger(new PeriodicTrigger(0)); - consumer.setPollerMetadata(pollerMetadata); + consumer.setTrigger(new PeriodicTrigger(0)); consumer.setErrorHandler(errorHandler); consumer.setTaskScheduler(taskScheduler); consumer.setBeanFactory(mock(BeanFactory.class)); @@ -111,7 +109,7 @@ public class PollingLifecycleTests { adapter.setTaskScheduler(taskScheduler); adapter.afterPropertiesSet(); adapter.start(); - assertTrue(latch.await(2, TimeUnit.SECONDS)); + assertTrue(latch.await(20, TimeUnit.SECONDS)); assertNotNull(channel.receive(100)); adapter.stop(); assertNull(channel.receive(1000)); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/pollingEndpointErrorHandlingTests.xml b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/pollingEndpointErrorHandlingTests.xml index ec67ab78ed..c234c89207 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/pollingEndpointErrorHandlingTests.xml +++ b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/pollingEndpointErrorHandlingTests.xml @@ -12,11 +12,7 @@ - - - + (new byte[] {1,2,3}), 0); channel.send(new GenericMessage(new byte[] {4,5,6}), 0); channel.send(new GenericMessage(new byte[] {7,8,9}), 0); @@ -116,10 +112,8 @@ public class ByteStreamWritingMessageHandlerTests { @Test public void maxMessagesPerTaskLessThanMessageCount() { - PollerMetadata pollerMetadata = new PollerMetadata(); - pollerMetadata.setTrigger(trigger); - pollerMetadata.setMaxMessagesPerPoll(2); - endpoint.setPollerMetadata(pollerMetadata); + endpoint.setTrigger(trigger); + endpoint.setMaxMessagesPerPoll(2); channel.send(new GenericMessage(new byte[] {1,2,3}), 0); channel.send(new GenericMessage(new byte[] {4,5,6}), 0); channel.send(new GenericMessage(new byte[] {7,8,9}), 0); @@ -133,10 +127,8 @@ public class ByteStreamWritingMessageHandlerTests { @Test public void maxMessagesPerTaskExceedsMessageCount() { - PollerMetadata pollerMetadata = new PollerMetadata(); - pollerMetadata.setTrigger(trigger); - pollerMetadata.setMaxMessagesPerPoll(5); - endpoint.setPollerMetadata(pollerMetadata); + endpoint.setTrigger(trigger); + endpoint.setMaxMessagesPerPoll(5); endpoint.setReceiveTimeout(0); channel.send(new GenericMessage(new byte[] {1,2,3}), 0); channel.send(new GenericMessage(new byte[] {4,5,6}), 0); @@ -151,10 +143,8 @@ public class ByteStreamWritingMessageHandlerTests { @Test public void testMaxMessagesLessThanMessageCountWithMultipleDispatches() { - PollerMetadata pollerMetadata = new PollerMetadata(); - pollerMetadata.setTrigger(trigger); - pollerMetadata.setMaxMessagesPerPoll(2); - endpoint.setPollerMetadata(pollerMetadata); + endpoint.setTrigger(trigger); + endpoint.setMaxMessagesPerPoll(2); endpoint.setReceiveTimeout(0); channel.send(new GenericMessage(new byte[] {1,2,3}), 0); channel.send(new GenericMessage(new byte[] {4,5,6}), 0); @@ -177,10 +167,8 @@ public class ByteStreamWritingMessageHandlerTests { @Test public void testMaxMessagesExceedsMessageCountWithMultipleDispatches() { - PollerMetadata pollerMetadata = new PollerMetadata(); - pollerMetadata.setTrigger(trigger); - pollerMetadata.setMaxMessagesPerPoll(5); - endpoint.setPollerMetadata(pollerMetadata); + endpoint.setTrigger(trigger); + endpoint.setMaxMessagesPerPoll(5); endpoint.setReceiveTimeout(0); channel.send(new GenericMessage(new byte[] {1,2,3}), 0); channel.send(new GenericMessage(new byte[] {4,5,6}), 0); @@ -202,10 +190,8 @@ public class ByteStreamWritingMessageHandlerTests { @Test public void testStreamResetBetweenDispatches() { - PollerMetadata pollerMetadata = new PollerMetadata(); - pollerMetadata.setMaxMessagesPerPoll(2); - pollerMetadata.setTrigger(trigger); - endpoint.setPollerMetadata(pollerMetadata); + endpoint.setMaxMessagesPerPoll(2); + endpoint.setTrigger(trigger); endpoint.setReceiveTimeout(0); channel.send(new GenericMessage(new byte[] {1,2,3}), 0); channel.send(new GenericMessage(new byte[] {4,5,6}), 0); @@ -227,10 +213,8 @@ public class ByteStreamWritingMessageHandlerTests { @Test public void testStreamWriteBetweenDispatches() throws IOException { - PollerMetadata pollerMetadata = new PollerMetadata(); - pollerMetadata.setTrigger(trigger); - pollerMetadata.setMaxMessagesPerPoll(2); - endpoint.setPollerMetadata(pollerMetadata); + endpoint.setTrigger(trigger); + endpoint.setMaxMessagesPerPoll(2); endpoint.setReceiveTimeout(0); channel.send(new GenericMessage(new byte[] {1,2,3}), 0); channel.send(new GenericMessage(new byte[] {4,5,6}), 0); diff --git a/spring-integration-stream/src/test/java/org/springframework/integration/stream/CharacterStreamWritingMessageHandlerTests.java b/spring-integration-stream/src/test/java/org/springframework/integration/stream/CharacterStreamWritingMessageHandlerTests.java index 7e2a0be5af..2a30ac4651 100644 --- a/spring-integration-stream/src/test/java/org/springframework/integration/stream/CharacterStreamWritingMessageHandlerTests.java +++ b/spring-integration-stream/src/test/java/org/springframework/integration/stream/CharacterStreamWritingMessageHandlerTests.java @@ -28,11 +28,11 @@ import java.util.concurrent.atomic.AtomicBoolean; import org.junit.After; import org.junit.Before; import org.junit.Test; + import org.springframework.beans.factory.BeanFactory; import org.springframework.integration.channel.QueueChannel; import org.springframework.integration.endpoint.PollingConsumer; import org.springframework.integration.message.GenericMessage; -import org.springframework.integration.scheduling.PollerMetadata; import org.springframework.scheduling.Trigger; import org.springframework.scheduling.TriggerContext; import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler; @@ -66,9 +66,7 @@ public class CharacterStreamWritingMessageHandlerTests { this.endpoint.setTaskScheduler(scheduler); scheduler.afterPropertiesSet(); trigger.reset(); - PollerMetadata pollerMetadata = new PollerMetadata(); - pollerMetadata.setTrigger(trigger); - endpoint.setPollerMetadata(pollerMetadata); + endpoint.setTrigger(trigger); endpoint.setBeanFactory(mock(BeanFactory.class)); } @@ -86,10 +84,8 @@ public class CharacterStreamWritingMessageHandlerTests { @Test public void twoStringsAndNoNewLinesByDefault() { - PollerMetadata pollerMetadata = new PollerMetadata(); - pollerMetadata.setMaxMessagesPerPoll(1); - pollerMetadata.setTrigger(trigger); - endpoint.setPollerMetadata(pollerMetadata); + endpoint.setMaxMessagesPerPoll(1); + endpoint.setTrigger(trigger); channel.send(new GenericMessage("foo"), 0); channel.send(new GenericMessage("bar"), 0); endpoint.start(); @@ -106,10 +102,8 @@ public class CharacterStreamWritingMessageHandlerTests { @Test public void twoStringsWithNewLines() { handler.setShouldAppendNewLine(true); - PollerMetadata pollerMetadata = new PollerMetadata(); - pollerMetadata.setTrigger(trigger); - pollerMetadata.setMaxMessagesPerPoll(1); - endpoint.setPollerMetadata(pollerMetadata); + endpoint.setTrigger(trigger); + endpoint.setMaxMessagesPerPoll(1); channel.send(new GenericMessage("foo"), 0); channel.send(new GenericMessage("bar"), 0); endpoint.start(); @@ -126,10 +120,8 @@ public class CharacterStreamWritingMessageHandlerTests { @Test public void maxMessagesPerTaskSameAsMessageCount() { - PollerMetadata pollerMetadata = new PollerMetadata(); - pollerMetadata.setTrigger(trigger); - pollerMetadata.setMaxMessagesPerPoll(2); - endpoint.setPollerMetadata(pollerMetadata); + endpoint.setTrigger(trigger); + endpoint.setMaxMessagesPerPoll(2); channel.send(new GenericMessage("foo"), 0); channel.send(new GenericMessage("bar"), 0); endpoint.start(); @@ -140,10 +132,8 @@ public class CharacterStreamWritingMessageHandlerTests { @Test public void maxMessagesPerTaskExceedsMessageCountWithAppendedNewLines() { - PollerMetadata pollerMetadata = new PollerMetadata(); - pollerMetadata.setTrigger(trigger); - pollerMetadata.setMaxMessagesPerPoll(10); - endpoint.setPollerMetadata(pollerMetadata); + endpoint.setTrigger(trigger); + endpoint.setMaxMessagesPerPoll(10); endpoint.setReceiveTimeout(0); handler.setShouldAppendNewLine(true); channel.send(new GenericMessage("foo"), 0); @@ -157,10 +147,8 @@ public class CharacterStreamWritingMessageHandlerTests { @Test public void singleNonStringObject() { - PollerMetadata pollerMetadata = new PollerMetadata(); - pollerMetadata.setTrigger(trigger); - pollerMetadata.setMaxMessagesPerPoll(1); - endpoint.setPollerMetadata(pollerMetadata); + endpoint.setTrigger(trigger); + endpoint.setMaxMessagesPerPoll(1); TestObject testObject = new TestObject("foo"); channel.send(new GenericMessage(testObject)); endpoint.start(); @@ -172,10 +160,8 @@ public class CharacterStreamWritingMessageHandlerTests { @Test public void twoNonStringObjectWithOutNewLines() { endpoint.setReceiveTimeout(0); - PollerMetadata pollerMetadata = new PollerMetadata(); - pollerMetadata.setTrigger(trigger); - pollerMetadata.setMaxMessagesPerPoll(2); - endpoint.setPollerMetadata(pollerMetadata); + endpoint.setTrigger(trigger); + endpoint.setMaxMessagesPerPoll(2); TestObject testObject1 = new TestObject("foo"); TestObject testObject2 = new TestObject("bar"); channel.send(new GenericMessage(testObject1), 0); @@ -190,10 +176,8 @@ public class CharacterStreamWritingMessageHandlerTests { public void twoNonStringObjectWithNewLines() { handler.setShouldAppendNewLine(true); endpoint.setReceiveTimeout(0); - PollerMetadata pollerMetadata = new PollerMetadata(); - pollerMetadata.setMaxMessagesPerPoll(2); - pollerMetadata.setTrigger(trigger); - endpoint.setPollerMetadata(pollerMetadata); + endpoint.setMaxMessagesPerPoll(2); + endpoint.setTrigger(trigger); TestObject testObject1 = new TestObject("foo"); TestObject testObject2 = new TestObject("bar"); channel.send(new GenericMessage(testObject1), 0); diff --git a/spring-integration-ws/src/test/java/org/springframework/integration/ws/config/WebServiceOutboundGatewayParserTests.java b/spring-integration-ws/src/test/java/org/springframework/integration/ws/config/WebServiceOutboundGatewayParserTests.java index 80db088230..61a3461cf4 100644 --- a/spring-integration-ws/src/test/java/org/springframework/integration/ws/config/WebServiceOutboundGatewayParserTests.java +++ b/spring-integration-ws/src/test/java/org/springframework/integration/ws/config/WebServiceOutboundGatewayParserTests.java @@ -229,9 +229,7 @@ public class WebServiceOutboundGatewayParserTests { "simpleWebServiceOutboundGatewayParserTests.xml", this.getClass()); AbstractEndpoint endpoint = (AbstractEndpoint) context.getBean("gatewayWithPoller"); assertEquals(PollingConsumer.class, endpoint.getClass()); - Object pollerMetadata = new DirectFieldAccessor(endpoint).getPropertyValue("pollerMetadata"); - assertEquals(PollerMetadata.class, pollerMetadata.getClass()); - Object triggerObject = new DirectFieldAccessor(pollerMetadata).getPropertyValue("trigger"); + Object triggerObject = new DirectFieldAccessor(endpoint).getPropertyValue("trigger"); assertEquals(PeriodicTrigger.class, triggerObject.getClass()); PeriodicTrigger trigger = (PeriodicTrigger) triggerObject; DirectFieldAccessor accessor = new DirectFieldAccessor(trigger);