INT-1923 refactored Lifecycle support for IMAP adapter that was fixed earlier
This commit is contained in:
@@ -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");
|
||||
|
||||
@@ -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");
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user