diff --git a/org.springframework.integration.mail/src/main/java/org/springframework/integration/mail/DefaultFolderConnection.java b/org.springframework.integration.mail/src/main/java/org/springframework/integration/mail/DefaultFolderConnection.java index 14313fc677..8e3df7930d 100644 --- a/org.springframework.integration.mail/src/main/java/org/springframework/integration/mail/DefaultFolderConnection.java +++ b/org.springframework.integration.mail/src/main/java/org/springframework/integration/mail/DefaultFolderConnection.java @@ -66,11 +66,11 @@ public class DefaultFolderConnection implements Lifecycle, DisposableBean, Folde public DefaultFolderConnection(String storeUri, MonitoringStrategy monitoringStrategy, boolean polling) { + Assert.notNull(storeUri, "storeUri must not be null"); + Assert.notNull(monitoringStrategy, "monitoringStrategy must not ne null"); this.storeUri = new URLName(storeUri); this.monitoringStrategy = monitoringStrategy; this.polling = polling; - Assert.notNull(storeUri, "storeUri is required"); - Assert.notNull(monitoringStrategy, "monitoringStrategy is required"); if (!polling && monitoringStrategy.getClass().isAssignableFrom(AsyncMonitoringStrategy.class)) { throw new ConfigurationException( "Folder connection requires an AsyncMonitoringStrategy if polling is disabled."); @@ -96,6 +96,11 @@ public class DefaultFolderConnection implements Lifecycle, DisposableBean, Folde } } + @Override + public String toString() { + return this.storeUri.toString(); + } + public void destroy() throws Exception { this.stop(); } diff --git a/org.springframework.integration.mail/src/main/java/org/springframework/integration/mail/ListeningMailSource.java b/org.springframework.integration.mail/src/main/java/org/springframework/integration/mail/ListeningMailSource.java index 42cb9f1299..97f029df69 100644 --- a/org.springframework.integration.mail/src/main/java/org/springframework/integration/mail/ListeningMailSource.java +++ b/org.springframework.integration.mail/src/main/java/org/springframework/integration/mail/ListeningMailSource.java @@ -16,8 +16,6 @@ package org.springframework.integration.mail; -import java.util.Date; - import javax.mail.Message; import org.apache.commons.logging.Log; @@ -25,10 +23,11 @@ import org.apache.commons.logging.LogFactory; import org.springframework.beans.factory.DisposableBean; import org.springframework.context.Lifecycle; +import org.springframework.core.task.SimpleAsyncTaskExecutor; +import org.springframework.core.task.TaskExecutor; import org.springframework.integration.endpoint.AbstractMessageProducingEndpoint; import org.springframework.integration.mail.monitor.AsyncMonitoringStrategy; import org.springframework.integration.message.MessageBuilder; -import org.springframework.integration.scheduling.Trigger; import org.springframework.util.Assert; /** @@ -44,6 +43,8 @@ public class ListeningMailSource extends AbstractMessageProducingEndpoint implem private final Log logger = LogFactory.getLog(this.getClass()); + private volatile TaskExecutor taskExecutor; + private final MonitorRunnable monitorRunnable; private volatile boolean monitorRunning = false; @@ -64,28 +65,30 @@ public class ListeningMailSource extends AbstractMessageProducingEndpoint implem public void start() { this.startMonitor(); if (logger.isInfoEnabled()) { - logger.info("started monitoring mailbox"); + logger.info("started monitoring mailbox [" + + this.monitorRunnable.folderConnection + "]"); } } public void stop() { this.stopMonitor(); if (logger.isInfoEnabled()) { - logger.info("stopped monitoring mailbox"); + logger.info("stopped monitoring mailbox [" + + this.monitorRunnable.folderConnection + "]"); } } protected void startMonitor() { synchronized (this.monitorRunnable) { if (!this.monitorRunning) { - Assert.state(this.getTaskScheduler() != null, "TaskScheduler is required"); - this.getTaskScheduler().schedule(this.monitorRunnable, - new Trigger() { - public Date getNextRunTime(Date lastScheduledRunTime, Date lastCompleteTime) { - return new Date(); - } - } - ); + if (this.taskExecutor == null) { + if (logger.isInfoEnabled()) { + logger.info("No TaskExecutor has been provided, will use a [" + + SimpleAsyncTaskExecutor.class + "] as the default."); + } + this.taskExecutor = new SimpleAsyncTaskExecutor(); + } + this.taskExecutor.execute(this.monitorRunnable); } this.monitorRunning = true; } @@ -113,6 +116,7 @@ public class ListeningMailSource extends AbstractMessageProducingEndpoint implem private MonitorRunnable(FolderConnection folderConnection) { + Assert.notNull(folderConnection, "folderConnection must not be null"); this.folderConnection = folderConnection; } diff --git a/org.springframework.integration.mail/src/test/java/org/springframework/integration/mail/SubscribableMailSourceTests.java b/org.springframework.integration.mail/src/test/java/org/springframework/integration/mail/SubscribableMailSourceTests.java index 246e6ab094..a6d8853937 100644 --- a/org.springframework.integration.mail/src/test/java/org/springframework/integration/mail/SubscribableMailSourceTests.java +++ b/org.springframework.integration.mail/src/test/java/org/springframework/integration/mail/SubscribableMailSourceTests.java @@ -26,11 +26,8 @@ import javax.mail.internet.MimeMessage; import org.easymock.classextension.EasyMock; import org.junit.Test; -import org.springframework.core.task.SimpleAsyncTaskExecutor; import org.springframework.integration.channel.QueueChannel; import org.springframework.integration.message.Message; -import org.springframework.integration.scheduling.SimpleTaskScheduler; -import org.springframework.integration.scheduling.TaskScheduler; /** * @author Jonas Partner @@ -42,10 +39,7 @@ public class SubscribableMailSourceTests { javax.mail.Message message = EasyMock.createMock(MimeMessage.class); StubFolderConnection folderConnection = new StubFolderConnection(message); QueueChannel channel = new QueueChannel(); - TaskScheduler scheduler = new SimpleTaskScheduler(new SimpleAsyncTaskExecutor()); - scheduler.start(); ListeningMailSource mailSource = new ListeningMailSource(folderConnection); - mailSource.setTaskScheduler(scheduler); mailSource.setOutputChannel(channel); mailSource.start(); Message result = channel.receive(1000); @@ -53,7 +47,6 @@ public class SubscribableMailSourceTests { assertNotNull(result); assertEquals("Wrong payload", message, result.getPayload()); mailSource.stop(); - scheduler.stop(); } diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/AbstractEndpoint.java b/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/AbstractEndpoint.java index 7f9dd7e05b..820a3a2cc6 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/AbstractEndpoint.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/AbstractEndpoint.java @@ -20,35 +20,23 @@ import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import org.springframework.beans.factory.BeanNameAware; -import org.springframework.integration.scheduling.TaskScheduler; -import org.springframework.integration.scheduling.TaskSchedulerAware; /** * The base class for Message Endpoint implementations. * * @author Mark Fisher */ -public abstract class AbstractEndpoint implements MessageEndpoint, TaskSchedulerAware, BeanNameAware { +public abstract class AbstractEndpoint implements MessageEndpoint, BeanNameAware { protected final Log logger = LogFactory.getLog(this.getClass()); private volatile String name; - private volatile TaskScheduler taskScheduler; - public void setBeanName(String name) { this.name = name; } - protected TaskScheduler getTaskScheduler() { - return this.taskScheduler; - } - - public void setTaskScheduler(TaskScheduler taskScheduler) { - this.taskScheduler = taskScheduler; - } - public String toString() { return (this.name != null) ? this.name : super.toString(); }