INT-1923 polished IDLE adapter to accomodate deprecation of task-executor
This commit is contained in:
@@ -26,6 +26,7 @@ import javax.mail.MessagingException;
|
||||
import javax.mail.Store;
|
||||
import javax.mail.internet.MimeMessage;
|
||||
|
||||
import org.springframework.core.task.SimpleAsyncTaskExecutor;
|
||||
import org.springframework.integration.endpoint.MessageProducerSupport;
|
||||
import org.springframework.integration.support.MessageBuilder;
|
||||
import org.springframework.scheduling.TaskScheduler;
|
||||
@@ -48,35 +49,41 @@ public class ImapIdleChannelAdapter extends MessageProducerSupport {
|
||||
|
||||
private volatile boolean shouldReconnectAutomatically = true;
|
||||
|
||||
private volatile Executor taskExecutor;
|
||||
private volatile Executor taskExecutor = new SimpleAsyncTaskExecutor();
|
||||
|
||||
private final ImapMailReceiver mailReceiver;
|
||||
|
||||
|
||||
private volatile int reconnectDelay = 10000; // seconds
|
||||
|
||||
|
||||
private volatile ScheduledFuture<?> receivingTask;
|
||||
private volatile ScheduledFuture<?> pingTask;
|
||||
|
||||
|
||||
private volatile long connectionPingInterval = 10000;
|
||||
|
||||
public ImapIdleChannelAdapter(ImapMailReceiver mailReceiver) {
|
||||
Assert.notNull(mailReceiver, "mailReceiver must not be null");
|
||||
this.mailReceiver = mailReceiver;
|
||||
}
|
||||
|
||||
/**
|
||||
* Specify whether the IDLE task should reconnect automatically after
|
||||
* catching a {@link FolderClosedException} while waiting for messages.
|
||||
* The default value is <code>true</code>.
|
||||
* catching a {@link FolderClosedException} while waiting for messages. The
|
||||
* default value is <code>true</code>.
|
||||
*/
|
||||
public void setShouldReconnectAutomatically(boolean shouldReconnectAutomatically) {
|
||||
public void setShouldReconnectAutomatically(
|
||||
boolean shouldReconnectAutomatically) {
|
||||
this.shouldReconnectAutomatically = shouldReconnectAutomatically;
|
||||
}
|
||||
|
||||
|
||||
/**
|
||||
* @deprecated As of releae 2.0.5
|
||||
* @param taskExecutor
|
||||
*/
|
||||
public void setTaskExecutor(Executor taskExecutor) {
|
||||
this.taskExecutor = taskExecutor;
|
||||
}
|
||||
|
||||
public String getComponentType(){
|
||||
|
||||
public String getComponentType() {
|
||||
return "mail:imap-idle-channel-adapter";
|
||||
}
|
||||
|
||||
@@ -85,6 +92,7 @@ public class ImapIdleChannelAdapter extends MessageProducerSupport {
|
||||
logger.warn("error occurred in idle task", e);
|
||||
}
|
||||
}
|
||||
|
||||
/*
|
||||
* Lifecycle implementation
|
||||
*/
|
||||
@@ -94,22 +102,30 @@ public class ImapIdleChannelAdapter extends MessageProducerSupport {
|
||||
final TaskScheduler scheduler = this.getTaskScheduler();
|
||||
Assert.notNull(scheduler, "'taskScheduler' must not be null" );
|
||||
|
||||
receivingTask = scheduler.scheduleAtFixedRate(new Runnable(){
|
||||
receivingTask = scheduler.schedule(new Runnable(){
|
||||
|
||||
public void run() {
|
||||
try {
|
||||
idleTask.run();
|
||||
if (logger.isDebugEnabled()){
|
||||
logger.debug("Task completed successfully. Re-scheduling it again right away");
|
||||
}
|
||||
}
|
||||
catch (IllegalStateException e) { //run again after a delay
|
||||
logger.warn("Failed to execute IDLE task. Will atempt to resubmit in " + reconnectDelay + " milliseconds", e);
|
||||
scheduler.scheduleAtFixedRate(this, new Date(System.currentTimeMillis() + reconnectDelay), 1);
|
||||
}
|
||||
taskExecutor.execute(new Runnable(){
|
||||
|
||||
public void run() {
|
||||
try {
|
||||
idleTask.run();
|
||||
if (mailReceiver.getFolder().isOpen()){
|
||||
if (logger.isDebugEnabled()){
|
||||
logger.debug("Task completed successfully. Re-scheduling it again right away");
|
||||
}
|
||||
scheduler.schedule(this, new Date());
|
||||
}
|
||||
}
|
||||
catch (IllegalStateException e) { //run again after a delay
|
||||
logger.warn("Failed to execute IDLE task. Will atempt to resubmit in " + reconnectDelay + " milliseconds", e);
|
||||
scheduler.schedule(this, new Date(System.currentTimeMillis() + reconnectDelay));
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
}, 1);
|
||||
}, new Date());
|
||||
|
||||
pingTask = scheduler.scheduleAtFixedRate(new Runnable() {
|
||||
public void run() {
|
||||
@@ -125,16 +141,18 @@ public class ImapIdleChannelAdapter extends MessageProducerSupport {
|
||||
}, connectionPingInterval);
|
||||
}
|
||||
|
||||
@Override // guarded by super#lifecycleLock
|
||||
@Override
|
||||
// guarded by super#lifecycleLock
|
||||
protected void doStop() {
|
||||
receivingTask.cancel(true);
|
||||
pingTask.cancel(true);
|
||||
pingTask.cancel(true);
|
||||
try {
|
||||
mailReceiver.destroy();
|
||||
} catch (Exception e) {
|
||||
throw new IllegalStateException("Failure during the destruction of " + mailReceiver, e);
|
||||
throw new IllegalStateException(
|
||||
"Failure during the destruction of " + mailReceiver, e);
|
||||
}
|
||||
|
||||
|
||||
}
|
||||
|
||||
private class IdleTask implements Runnable {
|
||||
@@ -145,32 +163,36 @@ public class ImapIdleChannelAdapter extends MessageProducerSupport {
|
||||
logger.debug("waiting for mail");
|
||||
}
|
||||
mailReceiver.waitForNewMessages();
|
||||
if (mailReceiver.getFolder().isOpen()){
|
||||
if (mailReceiver.getFolder().isOpen()) {
|
||||
Message[] mailMessages = mailReceiver.receive();
|
||||
if (logger.isDebugEnabled()) {
|
||||
logger.debug("received " + mailMessages.length + " mail messages");
|
||||
logger.debug("received " + mailMessages.length
|
||||
+ " mail messages");
|
||||
}
|
||||
for (Message mailMessage : mailMessages) {
|
||||
final MimeMessage copied = new MimeMessage((MimeMessage) mailMessage);
|
||||
if (taskExecutor != null){
|
||||
taskExecutor.execute(new Runnable() {
|
||||
final MimeMessage copied = new MimeMessage(
|
||||
(MimeMessage) mailMessage);
|
||||
if (taskExecutor != null) {
|
||||
taskExecutor.execute(new Runnable() {
|
||||
public void run() {
|
||||
sendMessage(MessageBuilder.withPayload(copied).build());
|
||||
sendMessage(MessageBuilder.withPayload(
|
||||
copied).build());
|
||||
}
|
||||
});
|
||||
}
|
||||
else {
|
||||
sendMessage(MessageBuilder.withPayload(copied).build());
|
||||
} else {
|
||||
sendMessage(MessageBuilder.withPayload(copied)
|
||||
.build());
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
} catch (MessagingException e) {
|
||||
ImapIdleChannelAdapter.this.handleMailMessagingException(e);
|
||||
if (shouldReconnectAutomatically){
|
||||
throw new IllegalStateException("Failure in 'idle' task. Will resubmit", e);
|
||||
}
|
||||
else {
|
||||
throw new org.springframework.integration.MessagingException("Failure in 'idle' task. Will NOT resubmit", e);
|
||||
if (shouldReconnectAutomatically) {
|
||||
throw new IllegalStateException(
|
||||
"Failure in 'idle' task. Will resubmit", e);
|
||||
} else {
|
||||
throw new org.springframework.integration.MessagingException(
|
||||
"Failure in 'idle' task. Will NOT resubmit", e);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -107,7 +107,7 @@
|
||||
<xsd:attribute name="task-executor" type="xsd:string">
|
||||
<xsd:annotation>
|
||||
<xsd:documentation>
|
||||
Reference to a TaskExecutor which will be used to send inbound Mail to a MessageChannel.
|
||||
[DEPRECATED as of 2.0.5]
|
||||
</xsd:documentation>
|
||||
<xsd:appinfo>
|
||||
<tool:annotation kind="ref">
|
||||
|
||||
@@ -54,7 +54,6 @@ public class ImapIdleChannelAdapterParserTests {
|
||||
DirectFieldAccessor adapterAccessor = new DirectFieldAccessor(adapter);
|
||||
Object channel = context.getBean("channel");
|
||||
assertSame(channel, adapterAccessor.getPropertyValue("outputChannel"));
|
||||
assertNull(adapterAccessor.getPropertyValue("taskExecutor"));
|
||||
assertEquals(Boolean.FALSE, adapterAccessor.getPropertyValue("autoStartup"));
|
||||
Object receiver = adapterAccessor.getPropertyValue("mailReceiver");
|
||||
assertEquals(ImapMailReceiver.class, receiver.getClass());
|
||||
@@ -74,7 +73,6 @@ public class ImapIdleChannelAdapterParserTests {
|
||||
DirectFieldAccessor adapterAccessor = new DirectFieldAccessor(adapter);
|
||||
Object channel = context.getBean("channel");
|
||||
assertSame(channel, adapterAccessor.getPropertyValue("outputChannel"));
|
||||
assertNull(adapterAccessor.getPropertyValue("taskExecutor"));
|
||||
assertEquals(Boolean.FALSE, adapterAccessor.getPropertyValue("autoStartup"));
|
||||
Object receiver = adapterAccessor.getPropertyValue("mailReceiver");
|
||||
assertEquals(ImapMailReceiver.class, receiver.getClass());
|
||||
@@ -94,7 +92,6 @@ public class ImapIdleChannelAdapterParserTests {
|
||||
DirectFieldAccessor adapterAccessor = new DirectFieldAccessor(adapter);
|
||||
Object channel = context.getBean("channel");
|
||||
assertSame(channel, adapterAccessor.getPropertyValue("outputChannel"));
|
||||
assertNull(adapterAccessor.getPropertyValue("taskExecutor"));
|
||||
assertEquals(Boolean.FALSE, adapterAccessor.getPropertyValue("autoStartup"));
|
||||
Object receiver = adapterAccessor.getPropertyValue("mailReceiver");
|
||||
assertEquals(ImapMailReceiver.class, receiver.getClass());
|
||||
@@ -114,7 +111,6 @@ public class ImapIdleChannelAdapterParserTests {
|
||||
DirectFieldAccessor adapterAccessor = new DirectFieldAccessor(adapter);
|
||||
Object channel = context.getBean("channel");
|
||||
assertSame(channel, adapterAccessor.getPropertyValue("outputChannel"));
|
||||
assertNull(adapterAccessor.getPropertyValue("taskExecutor"));
|
||||
assertEquals(Boolean.FALSE, adapterAccessor.getPropertyValue("autoStartup"));
|
||||
Object receiver = adapterAccessor.getPropertyValue("mailReceiver");
|
||||
assertEquals(ImapMailReceiver.class, receiver.getClass());
|
||||
|
||||
Reference in New Issue
Block a user