diff --git a/spring-integration-mail/src/main/java/org/springframework/integration/mail/AbstractMailReceiver.java b/spring-integration-mail/src/main/java/org/springframework/integration/mail/AbstractMailReceiver.java index c0cba026c9..ed3fae4bd0 100755 --- a/spring-integration-mail/src/main/java/org/springframework/integration/mail/AbstractMailReceiver.java +++ b/spring-integration-mail/src/main/java/org/springframework/integration/mail/AbstractMailReceiver.java @@ -349,6 +349,7 @@ public abstract class AbstractMailReceiver extends IntegrationObjectSupport impl protected void onInit() throws Exception { super.onInit(); this.folderOpenMode = Folder.READ_WRITE; + this.initialized = true; } @Override 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 55533a399b..b571669e8a 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 @@ -23,6 +23,7 @@ import java.util.concurrent.ScheduledFuture; import javax.mail.FolderClosedException; import javax.mail.Message; import javax.mail.MessagingException; +import javax.mail.Store; import javax.mail.internet.MimeMessage; import org.springframework.integration.endpoint.MessageProducerSupport; @@ -54,7 +55,10 @@ public class ImapIdleChannelAdapter extends MessageProducerSupport { private volatile int reconnectDelay = 10000; // seconds private volatile ScheduledFuture scheduledFuture; - + + private volatile ResubmittingTask task; + + private volatile long connectionPingInterval = 10000; public ImapIdleChannelAdapter(ImapMailReceiver mailReceiver) { Assert.notNull(mailReceiver, "mailReceiver must not be null"); @@ -90,14 +94,37 @@ public class ImapIdleChannelAdapter extends MessageProducerSupport { protected void doStart() { TaskScheduler scheduler = this.getTaskScheduler(); Assert.notNull(scheduler, "'taskScheduler' must not be null" ); - ResubmittingTask task = new ResubmittingTask(this.idleTask, scheduler, reconnectDelay); + 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) { + } + } + }, connectionPingInterval); + } } @Override // guarded by super#lifecycleLock protected void doStop() { scheduledFuture.cancel(true); + this.task.stop(); + try { + mailReceiver.destroy(); + } catch (Exception e) { + throw new IllegalStateException("Failure during the destruction of " + mailReceiver, e); + } + } private class IdleTask implements Runnable { @@ -108,14 +135,16 @@ public class ImapIdleChannelAdapter extends MessageProducerSupport { logger.debug("waiting for mail"); } mailReceiver.waitForNewMessages(); - Message[] mailMessages = mailReceiver.receive(); - if (logger.isDebugEnabled()) { - logger.debug("received " + mailMessages.length + " mail messages"); - } - for (Message mailMessage : mailMessages) { - MimeMessage copied = new MimeMessage((MimeMessage) mailMessage); - sendMessage(MessageBuilder.withPayload(copied).build()); - } + if (task.isRunning()){ + Message[] mailMessages = mailReceiver.receive(); + if (logger.isDebugEnabled()) { + logger.debug("received " + mailMessages.length + " mail messages"); + } + for (Message mailMessage : mailMessages) { + MimeMessage copied = new MimeMessage((MimeMessage) mailMessage); + sendMessage(MessageBuilder.withPayload(copied).build()); + } + } } catch (MessagingException e) { ImapIdleChannelAdapter.this.handleMailMessagingException(e); if (shouldReconnectAutomatically){ diff --git a/spring-integration-mail/src/main/java/org/springframework/integration/mail/ImapMailReceiver.java b/spring-integration-mail/src/main/java/org/springframework/integration/mail/ImapMailReceiver.java index c9f8247134..05b9799c42 100755 --- a/spring-integration-mail/src/main/java/org/springframework/integration/mail/ImapMailReceiver.java +++ b/spring-integration-mail/src/main/java/org/springframework/integration/mail/ImapMailReceiver.java @@ -21,7 +21,6 @@ import javax.mail.Flags.Flag; import javax.mail.Folder; import javax.mail.Message; import javax.mail.MessagingException; -import javax.mail.Store; import javax.mail.event.MessageCountAdapter; import javax.mail.event.MessageCountEvent; import javax.mail.event.MessageCountListener; @@ -30,7 +29,6 @@ import javax.mail.search.FlagTerm; import javax.mail.search.NotTerm; import javax.mail.search.SearchTerm; -import org.springframework.scheduling.TaskScheduler; import org.springframework.util.Assert; import com.sun.mail.imap.IMAPFolder; @@ -54,9 +52,6 @@ public class ImapMailReceiver extends AbstractMailReceiver { private final MessageCountListener messageCountListener = new SimpleMessageCountListener(); - private volatile long connectionPingInterval = 10000; - - public ImapMailReceiver() { super(); this.setProtocol("imap"); @@ -196,27 +191,6 @@ public class ImapMailReceiver extends AbstractMailReceiver { return searchTerm; } - @Override - protected void onInit() throws Exception { - super.onInit(); - this.initialized = true; - TaskScheduler scheduler = this.getTaskScheduler(); - if (scheduler != null) { - scheduler.scheduleAtFixedRate(new Runnable() { - public void run() { - try { - Store store = getStore(); - if (initialized && store != null) { - store.isConnected(); - } - } - catch (Exception ignore) { - } - } - }, connectionPingInterval); - } - } - protected void setAdditionalFlags(Message message) throws MessagingException { super.setAdditionalFlags(message); if (this.shouldMarkMessagesAsRead) { 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 29dc7a25ba..9e5cf01f8a 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,13 +35,15 @@ 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; + public ResubmittingTask(Runnable targetTask, TaskScheduler scheduler, long delay) { this.targetTask = targetTask; this.scheduler = scheduler; @@ -60,15 +62,37 @@ class ResubmittingTask implements Runnable { }); } + protected void stop(){ + this.running = false; + } + + protected void start(){ + this.running = true; + } + + protected boolean isRunning(){ + return this.running; + } + private void invokeTask(){ try { targetTask.run(); - logger.debug("Task completed successfully. Re-scheduling it again right away"); - scheduler.schedule(this, new Date()); + if (this.running){ + if (logger.isDebugEnabled()){ + logger.debug("Task completed successfully. Re-scheduling it again right away"); + } + scheduler.schedule(this, new Date()); + } + else { + if (logger.isDebugEnabled()){ + logger.debug("IDLE Task is stopped"); + } + } + } catch (IllegalStateException e) { //run again after a delay logger.warn("Failed to execute IDLE task. Will atempt to resubmit in " + delay + " milliseconds", e); - scheduler.schedule(this, new Date(System.currentTimeMillis() + delay)); + scheduler.schedule(this, new Date(System.currentTimeMillis() + delay)); } } }