AbstractEndpoint is no longer TaskScheduler aware. The ListeningMailSource now provides its own taskExecutor property and uses a SimpleAsyncTaskExecutor by default.
This commit is contained in:
@@ -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();
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
|
||||
@@ -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();
|
||||
}
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user