From 33e52486cf4d4f75a0105dd71e78db607bc515ba Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Wed, 12 Dec 2018 12:14:11 -0500 Subject: [PATCH] More fixes for IMAP channel adapter and its tests https://build.spring.io/browse/INT-MASTERSPRING40-553 The `ImapMailReceiverTests` fails sporadically according some race condition or wrong logic. * Fix `ImapIdleChannelAdapter` to check for folder not null before performing logic in the `IdleTask` * Polishing for logs which are based on constant strings * Remove `volatile` from configuration properties in the `ImapIdleChannelAdapter` and `AbstractMailReceiver` * Refactor some smells into the `protected` getters instead of direct access to the property * Stop channel adapters in the `ImapMailReceiverTests` * Also destroy task schedulers in the `ImapMailReceiverTests` --- .../mail/AbstractMailReceiver.java | 76 ++++++++++--------- .../mail/ImapIdleChannelAdapter.java | 64 +++++++--------- .../mail/ImapMailReceiverTests.java | 25 ++++-- 3 files changed, 86 insertions(+), 79 deletions(-) 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 3204f7a1e3..4c77abf4c7 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 @@ -60,6 +60,7 @@ import org.springframework.util.FileCopyUtils; * @author Iwein Fuld * @author Oleg Zhurakousky * @author Gary Russell + * @author Artem Bilan */ public abstract class AbstractMailReceiver extends IntegrationObjectSupport implements MailReceiver, DisposableBean { @@ -69,44 +70,42 @@ public abstract class AbstractMailReceiver extends IntegrationObjectSupport impl */ public final static String DEFAULT_SI_USER_FLAG = "spring-integration-mail-adapter"; - protected final Log logger = LogFactory.getLog(getClass()); + protected final Log logger = LogFactory.getLog(getClass()); // NOSONAR safe to use final private final URLName url; private final Object folderMonitor = new Object(); - private volatile String protocol; + private String protocol; - private volatile int maxFetchSize = -1; + private int maxFetchSize = -1; - private volatile Session session; + private Session session; + + private boolean shouldDeleteMessages; + + private int folderOpenMode = Folder.READ_ONLY; + + private Properties javaMailProperties = new Properties(); + + private Authenticator javaMailAuthenticator; + + private StandardEvaluationContext evaluationContext; + + private Expression selectorExpression; + + private HeaderMapper headerMapper; + + private String userFlag = DEFAULT_SI_USER_FLAG; + + private boolean embeddedPartsAsBytes = true; + + private boolean simpleContent; private volatile Store store; private volatile Folder folder; - private volatile boolean shouldDeleteMessages; - - protected volatile int folderOpenMode = Folder.READ_ONLY; - - private volatile Properties javaMailProperties = new Properties(); - - private volatile Authenticator javaMailAuthenticator; - - private volatile StandardEvaluationContext evaluationContext; - - private volatile Expression selectorExpression; - - private volatile HeaderMapper headerMapper; - - protected volatile boolean initialized; - - private volatile String userFlag = DEFAULT_SI_USER_FLAG; - - private volatile boolean embeddedPartsAsBytes = true; - - private volatile boolean simpleContent; - public AbstractMailReceiver() { this.url = null; } @@ -283,6 +282,10 @@ public abstract class AbstractMailReceiver extends IntegrationObjectSupport impl return this.folder; } + protected int getFolderOpenMode() { + return this.folderOpenMode; + } + /** * Subclasses must implement this method to return new mail messages. * @@ -291,7 +294,7 @@ public abstract class AbstractMailReceiver extends IntegrationObjectSupport impl */ protected abstract Message[] searchForNewMessages() throws MessagingException; - private void openSession() throws MessagingException { + private void openSession() { if (this.session == null) { if (this.javaMailAuthenticator != null) { this.session = Session.getInstance(this.javaMailProperties, this.javaMailAuthenticator); @@ -316,7 +319,7 @@ public abstract class AbstractMailReceiver extends IntegrationObjectSupport impl } if (!this.store.isConnected()) { if (this.logger.isDebugEnabled()) { - this.logger.debug("connecting to store [" + MailTransportUtils.toPasswordProtectedString(this.url) + "]"); + this.logger.debug("connecting to store [" + this.store.getURLName() + "]"); } this.store.connect(); } @@ -338,7 +341,7 @@ public abstract class AbstractMailReceiver extends IntegrationObjectSupport impl return; } if (this.logger.isDebugEnabled()) { - this.logger.debug("opening folder [" + MailTransportUtils.toPasswordProtectedString(this.url) + "]"); + this.logger.debug("opening folder [" + this.folder.getURLName() + "]"); } this.folder.open(this.folderOpenMode); } @@ -353,7 +356,7 @@ public abstract class AbstractMailReceiver extends IntegrationObjectSupport impl try { this.openFolder(); if (this.logger.isInfoEnabled()) { - this.logger.info("attempting to receive mail from folder [" + this.getFolder().getFullName() + "]"); + this.logger.info("attempting to receive mail from folder [" + getFolder().getFullName() + "]"); } Message[] messages = searchForNewMessages(); if (this.maxFetchSize > 0 && messages.length > this.maxFetchSize) { @@ -496,17 +499,18 @@ public abstract class AbstractMailReceiver extends IntegrationObjectSupport impl * will be filtered out and remain on the server as never touched. */ private MimeMessage[] filterMessagesThruSelector(Message[] messages) throws MessagingException { - List filteredMessages = new LinkedList(); + List filteredMessages = new LinkedList<>(); for (int i = 0; i < messages.length; i++) { MimeMessage message = (MimeMessage) messages[i]; if (this.selectorExpression != null) { - if (this.selectorExpression.getValue(this.evaluationContext, message, Boolean.class)) { + if (Boolean.TRUE.equals( + this.selectorExpression.getValue(this.evaluationContext, message, Boolean.class))) { filteredMessages.add(message); } else { if (this.logger.isDebugEnabled()) { - this.logger.debug("Fetched email with subject '" + message.getSubject() + "' will be discarded by the matching filter" + - " and will not be flagged as SEEN."); + this.logger.debug("Fetched email with subject '" + message.getSubject() + + "' will be discarded by the matching filter and will not be flagged as SEEN."); } } } @@ -514,7 +518,7 @@ public abstract class AbstractMailReceiver extends IntegrationObjectSupport impl filteredMessages.add(message); } } - return filteredMessages.toArray(new MimeMessage[filteredMessages.size()]); + return filteredMessages.toArray(new MimeMessage[0]); } /** @@ -562,7 +566,6 @@ public abstract class AbstractMailReceiver extends IntegrationObjectSupport impl MailTransportUtils.closeService(this.store); this.folder = null; this.store = null; - this.initialized = false; } } @@ -571,7 +574,6 @@ public abstract class AbstractMailReceiver extends IntegrationObjectSupport impl super.onInit(); this.folderOpenMode = Folder.READ_WRITE; this.evaluationContext = ExpressionUtils.createStandardEvaluationContext(getBeanFactory()); - 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 a14924c1ac..bdd107a07c 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 @@ -1,5 +1,5 @@ /* - * Copyright 2002-2017 the original author or authors. + * Copyright 2002-2018 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -23,6 +23,7 @@ import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledFuture; +import javax.mail.Folder; import javax.mail.FolderClosedException; import javax.mail.Message; import javax.mail.MessagingException; @@ -64,30 +65,30 @@ public class ImapIdleChannelAdapter extends MessageProducerSupport implements Be private static final int DEFAULT_RECONNECT_DELAY = 10000; + private final ExceptionAwarePeriodicTrigger receivingTaskTrigger = new ExceptionAwarePeriodicTrigger(); + private final IdleTask idleTask = new IdleTask(); - private volatile Executor sendingTaskExecutor; - - private volatile boolean sendingTaskExecutorSet; - - private volatile boolean shouldReconnectAutomatically = true; - - private volatile ClassLoader classLoader; - - private volatile List adviceChain; - private final ImapMailReceiver mailReceiver; - private volatile long reconnectDelay = DEFAULT_RECONNECT_DELAY; // milliseconds + private TransactionSynchronizationFactory transactionSynchronizationFactory; + + private ClassLoader classLoader; + + private ApplicationEventPublisher applicationEventPublisher; + + private boolean shouldReconnectAutomatically = true; + + private Executor sendingTaskExecutor; + + private boolean sendingTaskExecutorSet; + + private List adviceChain; + + private long reconnectDelay = DEFAULT_RECONNECT_DELAY; // milliseconds private volatile ScheduledFuture receivingTask; - private final ExceptionAwarePeriodicTrigger receivingTaskTrigger = new ExceptionAwarePeriodicTrigger(); - - private volatile TransactionSynchronizationFactory transactionSynchronizationFactory; - - private volatile ApplicationEventPublisher applicationEventPublisher; - public ImapIdleChannelAdapter(ImapMailReceiver mailReceiver) { Assert.notNull(mailReceiver, "'mailReceiver' must not be null"); this.mailReceiver = mailReceiver; @@ -187,7 +188,7 @@ public class ImapIdleChannelAdapter extends MessageProducerSupport implements Be @SuppressWarnings("unchecked") org.springframework.messaging.Message message = mailMessage instanceof Message - ? ImapIdleChannelAdapter.this.getMessageBuilderFactory().withPayload(mailMessage).build() + ? getMessageBuilderFactory().withPayload(mailMessage).build() : (org.springframework.messaging.Message) mailMessage; if (TransactionSynchronizationManager.isActualTransactionActive()) { @@ -251,9 +252,7 @@ public class ImapIdleChannelAdapter extends MessageProducerSupport implements Be public void run() { try { ImapIdleChannelAdapter.this.idleTask.run(); - if (logger.isDebugEnabled()) { - logger.debug("Task completed successfully. Re-scheduling it again right away."); - } + logger.debug("Task completed successfully. Re-scheduling it again right away."); } catch (Exception e) { //run again after a delay if (logger.isWarnEnabled()) { @@ -281,33 +280,28 @@ public class ImapIdleChannelAdapter extends MessageProducerSupport implements Be * The following shouldn't be necessary because doStart() will have ensured we have * one. But, just in case... */ - Assert.state(ImapIdleChannelAdapter.this.sendingTaskExecutor != null, "'sendingTaskExecutor' must not be null"); + Assert.state(ImapIdleChannelAdapter.this.sendingTaskExecutor != null, + "'sendingTaskExecutor' must not be null"); try { - if (logger.isDebugEnabled()) { - logger.debug("waiting for mail"); - } + logger.debug("waiting for mail"); ImapIdleChannelAdapter.this.mailReceiver.waitForNewMessages(); - if (ImapIdleChannelAdapter.this.mailReceiver.getFolder().isOpen()) { + Folder folder = ImapIdleChannelAdapter.this.mailReceiver.getFolder(); + if (folder != null && folder.isOpen()) { Object[] mailMessages = ImapIdleChannelAdapter.this.mailReceiver.receive(); if (logger.isDebugEnabled()) { logger.debug("received " + mailMessages.length + " mail messages"); } - for (final Object mailMessage : mailMessages) { - + for (Object mailMessage : mailMessages) { Runnable messageSendingTask = createMessageSendingTask(mailMessage); - ImapIdleChannelAdapter.this.sendingTaskExecutor.execute(messageSendingTask); } } } catch (MessagingException e) { - if (logger.isWarnEnabled()) { - logger.warn("error occurred in idle task", e); - } + logger.warn("error occurred in idle task", e); if (ImapIdleChannelAdapter.this.shouldReconnectAutomatically) { - throw new IllegalStateException( - "Failure in 'idle' task. Will resubmit.", e); + throw new IllegalStateException("Failure in 'idle' task. Will resubmit.", e); } else { throw new org.springframework.messaging.MessagingException( diff --git a/spring-integration-mail/src/test/java/org/springframework/integration/mail/ImapMailReceiverTests.java b/spring-integration-mail/src/test/java/org/springframework/integration/mail/ImapMailReceiverTests.java index 9f24f78cc7..7f536c2c1d 100644 --- a/spring-integration-mail/src/test/java/org/springframework/integration/mail/ImapMailReceiverTests.java +++ b/spring-integration-mail/src/test/java/org/springframework/integration/mail/ImapMailReceiverTests.java @@ -64,7 +64,6 @@ import javax.mail.internet.MimeMessage; import javax.mail.search.AndTerm; import javax.mail.search.FlagTerm; import javax.mail.search.FromTerm; -import javax.mail.search.SearchTerm; import org.apache.commons.logging.Log; import org.junit.AfterClass; @@ -254,6 +253,8 @@ public class ImapMailReceiverTests { assertNotNull(channel.receive(10000)); // new message after idle assertNull(channel.receive(100)); // no new message after second and third idle verify(logger).debug("Canceling IDLE"); + + adapter.stop(); taskScheduler.shutdown(); assertTrue(imapIdleServer.assertReceived("storeUserFlag")); } @@ -470,6 +471,7 @@ public class ImapMailReceiverTests { @Ignore public void testMessageHistory() throws Exception { ImapIdleChannelAdapter adapter = this.context.getBean("simpleAdapter", ImapIdleChannelAdapter.class); + adapter.setReconnectDelay(1); AbstractMailReceiver receiver = new ImapMailReceiver(); receiver = spy(receiver); @@ -526,6 +528,7 @@ public class ImapMailReceiverTests { adapter.setOutputChannel(channel); QueueChannel errorChannel = new QueueChannel(); adapter.setErrorChannel(errorChannel); + adapter.setReconnectDelay(1); AbstractMailReceiver receiver = new ImapMailReceiver(); receiver = spy(receiver); @@ -568,6 +571,7 @@ public class ImapMailReceiverTests { QueueChannel channel = new QueueChannel(); adapter.setOutputChannel(channel); + adapter.setReconnectDelay(1); ImapMailReceiver receiver = new ImapMailReceiver("imap:foo"); receiver = spy(receiver); @@ -706,10 +710,14 @@ public class ImapMailReceiverTests { ThreadPoolTaskScheduler taskScheduler = new ThreadPoolTaskScheduler(); taskScheduler.initialize(); adapter.setTaskScheduler(taskScheduler); + adapter.setReconnectDelay(1); adapter.start(); assertTrue(latch.await(10, TimeUnit.SECONDS)); assertThat(theEvent.get().toString(), endsWith("cause=java.lang.IllegalStateException: Failure in 'idle' task. Will resubmit.]")); + + adapter.stop(); + taskScheduler.destroy(); } @Test // see INT-1801 @@ -777,9 +785,9 @@ public class ImapMailReceiverTests { org.springframework.messaging.Message received = messages[0]; Object content = received.getPayload(); assertThat(content, instanceOf(byte[].class)); - assertThat((String) received.getHeaders().get(MailHeaders.CONTENT_TYPE), + assertThat(received.getHeaders().get(MailHeaders.CONTENT_TYPE), equalTo("multipart/mixed;\r\n boundary=\"------------040903000701040401040200\"")); - assertThat((String) received.getHeaders().get(MessageHeaders.CONTENT_TYPE), + assertThat(received.getHeaders().get(MessageHeaders.CONTENT_TYPE), equalTo("application/octet-stream")); } @@ -789,8 +797,8 @@ public class ImapMailReceiverTests { receiver.setHeaderMapper(new DefaultMailHeaderMapper()); receiver.setEmbeddedPartsAsBytes(false); testAttachmentsGuts(receiver); - org.springframework.messaging.Message[] messages = (org.springframework.messaging.Message[]) receiver - .receive(); + org.springframework.messaging.Message[] messages = + (org.springframework.messaging.Message[]) receiver.receive(); Object content = messages[0].getPayload(); assertThat(content, instanceOf(Multipart.class)); assertEquals("bar", ((Multipart) content).getBodyPart(0).getContent().toString().trim()); @@ -804,7 +812,7 @@ public class ImapMailReceiverTests { given(folder.isOpen()).willReturn(true); Message message = new MimeMessage(null, new ClassPathResource("test.mail").getInputStream()); - given(folder.search((SearchTerm) Mockito.any())).willReturn(new Message[] { message }); + given(folder.search(Mockito.any())).willReturn(new Message[] { message }); given(store.getFolder(Mockito.any(URLName.class))).willReturn(folder); given(folder.getPermanentFlags()).willReturn(new Flags(Flags.Flag.USER)); DirectFieldAccessor df = new DirectFieldAccessor(receiver); @@ -821,6 +829,7 @@ public class ImapMailReceiverTests { ThreadPoolTaskScheduler taskScheduler = new ThreadPoolTaskScheduler(); taskScheduler.initialize(); adapter.setTaskScheduler(taskScheduler); + adapter.setReconnectDelay(1); adapter.start(); ExecutorService exec = TestUtils.getPropertyValue(adapter, "sendingTaskExecutor", ExecutorService.class); adapter.stop(); @@ -908,7 +917,7 @@ public class ImapMailReceiverTests { ThreadPoolTaskScheduler taskScheduler = new ThreadPoolTaskScheduler(); taskScheduler.initialize(); adapter.setTaskScheduler(taskScheduler); - adapter.setReconnectDelay(50); + adapter.setReconnectDelay(1); adapter.afterPropertiesSet(); final CountDownLatch latch = new CountDownLatch(3); adapter.setApplicationEventPublisher(e -> { @@ -917,6 +926,8 @@ public class ImapMailReceiverTests { adapter.start(); assertTrue(latch.await(60, TimeUnit.SECONDS)); verify(store, atLeast(3)).connect(); + + adapter.stop(); taskScheduler.shutdown(); }