diff --git a/spring-integration-mail/src/main/java/org/springframework/integration/mail/ImapIdleChannelAdapter.java b/spring-integration-mail/src/main/java/org/springframework/integration/mail/ImapIdleChannelAdapter.java index b571669e8a..b811ecafd6 100755 --- a/spring-integration-mail/src/main/java/org/springframework/integration/mail/ImapIdleChannelAdapter.java +++ b/spring-integration-mail/src/main/java/org/springframework/integration/mail/ImapIdleChannelAdapter.java @@ -54,7 +54,8 @@ public class ImapIdleChannelAdapter extends MessageProducerSupport { private volatile int reconnectDelay = 10000; // seconds - private volatile ScheduledFuture scheduledFuture; + private volatile ScheduledFuture receivingTask; + private volatile ScheduledFuture pingTask; private volatile ResubmittingTask task; @@ -95,30 +96,29 @@ public class ImapIdleChannelAdapter extends MessageProducerSupport { TaskScheduler scheduler = this.getTaskScheduler(); Assert.notNull(scheduler, "'taskScheduler' must not be null" ); this.task = new ResubmittingTask(this.idleTask, scheduler, reconnectDelay); - this.task.start(); task.setTaskExecutor(taskExecutor); - scheduledFuture = scheduler.schedule(task, new Date()); - scheduler = this.getTaskScheduler(); - if (scheduler != null) { - scheduler.scheduleAtFixedRate(new Runnable() { - public void run() { - try { - Store store = mailReceiver.getStore(); - if (store != null) { - store.isConnected(); - } - } - catch (Exception ignore) { - } + + receivingTask = scheduler.schedule(task, new Date()); + + pingTask = scheduler.scheduleAtFixedRate(new Runnable() { + public void run() { + try { + Store store = mailReceiver.getStore(); + if (store != null) { + store.isConnected(); + } + } + catch (Exception ignore) { } - }, connectionPingInterval); - } + } + }, connectionPingInterval); } @Override // guarded by super#lifecycleLock protected void doStop() { - scheduledFuture.cancel(true); - this.task.stop(); + this.task.requestStop(); + receivingTask.cancel(true); + pingTask.cancel(true); try { mailReceiver.destroy(); } catch (Exception e) { @@ -135,7 +135,7 @@ public class ImapIdleChannelAdapter extends MessageProducerSupport { logger.debug("waiting for mail"); } mailReceiver.waitForNewMessages(); - if (task.isRunning()){ + if (!task.isStopRequested()){ Message[] mailMessages = mailReceiver.receive(); if (logger.isDebugEnabled()) { logger.debug("received " + mailMessages.length + " mail messages"); diff --git a/spring-integration-mail/src/main/java/org/springframework/integration/mail/ResubmittingTask.java b/spring-integration-mail/src/main/java/org/springframework/integration/mail/ResubmittingTask.java index 9e5cf01f8a..e979d711b9 100644 --- a/spring-integration-mail/src/main/java/org/springframework/integration/mail/ResubmittingTask.java +++ b/spring-integration-mail/src/main/java/org/springframework/integration/mail/ResubmittingTask.java @@ -35,14 +35,14 @@ import org.springframework.scheduling.TaskScheduler; * with reconnection logic * Currently only used to manage IDLE task of ImapIdleChannelAdapter */ -class ResubmittingTask implements Runnable{ +class ResubmittingTask implements Runnable { private static final Log logger = LogFactory.getLog(ResubmittingTask.class); private final Runnable targetTask; private final TaskScheduler scheduler; private final long delay; private Executor taskExecutor = new SimpleAsyncTaskExecutor(); - private volatile boolean running; + private volatile boolean stopRequested = false; public ResubmittingTask(Runnable targetTask, TaskScheduler scheduler, long delay) { this.targetTask = targetTask; @@ -62,22 +62,18 @@ class ResubmittingTask implements Runnable{ }); } - protected void stop(){ - this.running = false; + protected void requestStop(){ + this.stopRequested = true; } - - protected void start(){ - this.running = true; - } - - protected boolean isRunning(){ - return this.running; + + protected boolean isStopRequested(){ + return this.stopRequested; } private void invokeTask(){ try { targetTask.run(); - if (this.running){ + if (!this.stopRequested){ if (logger.isDebugEnabled()){ logger.debug("Task completed successfully. Re-scheduling it again right away"); }