Refactored all task scheduling and poller triggering to use functionality now included in the Spring 3.0 core, and removed Spring Integration specific code that is now handled by that corresponding code in the core.

This commit is contained in:
Mark Fisher
2009-08-04 01:25:53 +00:00
parent 7e6d4a6e27
commit 67b6ea048a
47 changed files with 399 additions and 1070 deletions

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2008 the original author or authors.
* Copyright 2002-2009 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.
@@ -29,13 +29,13 @@ import org.junit.After;
import org.junit.Before;
import org.junit.Test;
import org.springframework.core.task.SimpleAsyncTaskExecutor;
import org.springframework.integration.channel.QueueChannel;
import org.springframework.integration.endpoint.PollingConsumer;
import org.springframework.integration.message.GenericMessage;
import org.springframework.integration.message.StringMessage;
import org.springframework.integration.scheduling.SimpleTaskScheduler;
import org.springframework.integration.scheduling.Trigger;
import org.springframework.scheduling.Trigger;
import org.springframework.scheduling.TriggerContext;
import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler;
/**
* @author Mark Fisher
@@ -52,7 +52,7 @@ public class ByteStreamWritingMessageHandlerTests {
private TestTrigger trigger = new TestTrigger();
private SimpleTaskScheduler scheduler;
private ThreadPoolTaskScheduler scheduler;
@Before
@@ -61,9 +61,9 @@ public class ByteStreamWritingMessageHandlerTests {
handler = new ByteStreamWritingMessageHandler(stream);
this.channel = new QueueChannel(10);
this.endpoint = new PollingConsumer(channel, handler);
scheduler = new SimpleTaskScheduler(new SimpleAsyncTaskExecutor());
scheduler = new ThreadPoolTaskScheduler();
this.endpoint.setTaskScheduler(scheduler);
scheduler.start();
scheduler.afterPropertiesSet();
trigger.reset();
endpoint.setTrigger(trigger);
}
@@ -235,8 +235,7 @@ public class ByteStreamWritingMessageHandlerTests {
private volatile CountDownLatch latch = new CountDownLatch(1);
public Date getNextRunTime(Date lastScheduledRunTime, Date lastCompleteTime) {
public Date nextExecutionTime(TriggerContext triggerContext) {
if (!hasRun.getAndSet(true)) {
return new Date();
}

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2008 the original author or authors.
* Copyright 2002-2009 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.
@@ -28,13 +28,13 @@ import org.junit.After;
import org.junit.Before;
import org.junit.Test;
import org.springframework.core.task.SimpleAsyncTaskExecutor;
import org.springframework.integration.channel.QueueChannel;
import org.springframework.integration.endpoint.PollingConsumer;
import org.springframework.integration.message.GenericMessage;
import org.springframework.integration.message.StringMessage;
import org.springframework.integration.scheduling.SimpleTaskScheduler;
import org.springframework.integration.scheduling.Trigger;
import org.springframework.scheduling.Trigger;
import org.springframework.scheduling.TriggerContext;
import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler;
/**
* @author Mark Fisher
@@ -51,7 +51,7 @@ public class CharacterStreamWritingMessageHandlerTests {
private TestTrigger trigger = new TestTrigger();
private SimpleTaskScheduler scheduler;
private ThreadPoolTaskScheduler scheduler;
@Before
@@ -61,9 +61,9 @@ public class CharacterStreamWritingMessageHandlerTests {
this.channel = new QueueChannel(10);
trigger.reset();
this.endpoint = new PollingConsumer(channel, handler);
scheduler = new SimpleTaskScheduler(new SimpleAsyncTaskExecutor());
scheduler = new ThreadPoolTaskScheduler();
this.endpoint.setTaskScheduler(scheduler);
scheduler.start();
scheduler.afterPropertiesSet();
trigger.reset();
endpoint.setTrigger(trigger);
}
@@ -201,8 +201,7 @@ public class CharacterStreamWritingMessageHandlerTests {
private volatile CountDownLatch latch = new CountDownLatch(1);
public Date getNextRunTime(Date lastScheduledRunTime, Date lastCompleteTime) {
public Date nextExecutionTime(TriggerContext triggerContext) {
if (!hasRun.getAndSet(true)) {
return new Date();
}

View File

@@ -22,11 +22,14 @@ import static org.junit.Assert.assertNotNull;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.springframework.beans.DirectFieldAccessor;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.ApplicationContext;
import org.springframework.integration.channel.MessagePublishingErrorHandler;
import org.springframework.integration.channel.NullChannel;
import org.springframework.integration.channel.PublishSubscribeChannel;
import org.springframework.integration.scheduling.SimpleTaskScheduler;
import org.springframework.integration.context.IntegrationContextUtils;
import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
@@ -58,9 +61,12 @@ public class DefaultConfigurationTests {
@Test
public void verifyTaskScheduler() {
Object taskScheduler = context.getBean("taskScheduler");
assertNotNull(taskScheduler);
assertEquals(SimpleTaskScheduler.class, taskScheduler.getClass());
Object taskScheduler = context.getBean(IntegrationContextUtils.TASK_SCHEDULER_BEAN_NAME);
assertEquals(ThreadPoolTaskScheduler.class, taskScheduler.getClass());
Object errorHandler = new DirectFieldAccessor(taskScheduler).getPropertyValue("errorHandler");
assertEquals(MessagePublishingErrorHandler.class, errorHandler.getClass());
Object defaultErrorChannel = new DirectFieldAccessor(errorHandler).getPropertyValue("defaultErrorChannel");
assertEquals(context.getBean(IntegrationContextUtils.ERROR_CHANNEL_BEAN_NAME), defaultErrorChannel);
}
}