From 36eb5b1edeb0083995e1e08bd36225c8e92b85e6 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Thu, 14 Jul 2011 15:57:25 -0400 Subject: [PATCH] INT-1923 polished IDLE adapter to accomodate deprecation of task-executor --- .../mail/ImapIdleChannelAdapter.java | 104 +++++++++++------- .../config/spring-integration-mail-2.0.xsd | 2 +- .../ImapIdleChannelAdapterParserTests.java | 4 - 3 files changed, 64 insertions(+), 46 deletions(-) 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 a92c091164..b1164fd7a4 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 @@ -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 true. + * catching a {@link FolderClosedException} while waiting for messages. The + * default value is true. */ - 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); } } } diff --git a/spring-integration-mail/src/main/resources/org/springframework/integration/mail/config/spring-integration-mail-2.0.xsd b/spring-integration-mail/src/main/resources/org/springframework/integration/mail/config/spring-integration-mail-2.0.xsd index 2d8b0786ab..79f0844566 100644 --- a/spring-integration-mail/src/main/resources/org/springframework/integration/mail/config/spring-integration-mail-2.0.xsd +++ b/spring-integration-mail/src/main/resources/org/springframework/integration/mail/config/spring-integration-mail-2.0.xsd @@ -107,7 +107,7 @@ - Reference to a TaskExecutor which will be used to send inbound Mail to a MessageChannel. + [DEPRECATED as of 2.0.5] diff --git a/spring-integration-mail/src/test/java/org/springframework/integration/mail/config/ImapIdleChannelAdapterParserTests.java b/spring-integration-mail/src/test/java/org/springframework/integration/mail/config/ImapIdleChannelAdapterParserTests.java index 19e41a41ef..c6282474e7 100644 --- a/spring-integration-mail/src/test/java/org/springframework/integration/mail/config/ImapIdleChannelAdapterParserTests.java +++ b/spring-integration-mail/src/test/java/org/springframework/integration/mail/config/ImapIdleChannelAdapterParserTests.java @@ -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());