INT-1923 fixed Lifecycle.stop() support for IMAP adapter
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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){
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user