From 20eb2922ee66393aae0818bf1896692ff22772cc Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Tue, 12 Oct 2010 05:42:10 -0400 Subject: [PATCH] INT-1493, ensured that default 'maxMessagesPerPoll' for SPCA is 1 and -1 for PC, added tests validating that both SPCA and PC stops --- ...ourcePollingChannelAdapterFactoryBean.java | 3 + .../endpoint/AbstractPollingEndpoint.java | 2 +- .../endpoint/PollingLifecycleTests.java | 116 ++++++++++++++++++ 3 files changed, 120 insertions(+), 1 deletion(-) create mode 100644 spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollingLifecycleTests.java 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 957fbee0a5..74b34789d6 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 @@ -125,6 +125,9 @@ public class SourcePollingChannelAdapterFactoryBean implements FactoryBean("foo")); + + MessageHandler handler = Mockito.spy(new MessageHandler() { + public void handleMessage(Message message) throws MessagingException { + latch.countDown(); + } + }); + PollingConsumer consumer = new PollingConsumer(channel, handler); + PollerMetadata pollerMetadata = new PollerMetadata(); + pollerMetadata.setTrigger(new PeriodicTrigger(0)); + consumer.setPollerMetadata(pollerMetadata); + consumer.setErrorHandler(errorHandler); + consumer.setTaskScheduler(taskScheduler); + consumer.setBeanFactory(mock(BeanFactory.class)); + consumer.afterPropertiesSet(); + consumer.start(); + assertTrue(latch.await(2, TimeUnit.SECONDS)); + consumer.stop(); + for (int i = 0; i < 10; i++) { + channel.send(new GenericMessage("foo")); + } + Mockito.verify(handler, times(1)).handleMessage(Mockito.any(Message.class)); + } + + @Test + public void ensurePollerTaskStopsForAdapter() throws Exception{ + final CountDownLatch latch = new CountDownLatch(1); + QueueChannel channel = new QueueChannel(); + + SourcePollingChannelAdapterFactoryBean adapterFactory = new SourcePollingChannelAdapterFactoryBean(); + PollerMetadata pollerMetadata = new PollerMetadata(); + pollerMetadata.setMaxMessagesPerPoll(-1); // should be overriden in FB + pollerMetadata.setTrigger(new PeriodicTrigger(2000)); + adapterFactory.setPollerMetadata(pollerMetadata); + MessageSource source = spy(new MessageSource() { + public Message receive() { + latch.countDown(); + return new GenericMessage("hello"); + } + }); + adapterFactory.setSource(source); + adapterFactory.setOutputChannel(channel); + adapterFactory.setBeanFactory(mock(ConfigurableBeanFactory.class)); + SourcePollingChannelAdapter adapter = adapterFactory.getObject(); + adapter.setTaskScheduler(taskScheduler); + adapter.afterPropertiesSet(); + adapter.start(); + assertTrue(latch.await(2, TimeUnit.SECONDS)); + assertNotNull(channel.receive(100)); + adapter.stop(); + assertNull(channel.receive(1000)); + Mockito.verify(source, times(1)).receive(); + } +}