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`
This commit is contained in:
Artem Bilan
2018-12-12 12:14:11 -05:00
committed by Gary Russell
parent ca2e70f236
commit 33e52486cf
3 changed files with 86 additions and 79 deletions

View File

@@ -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<MimeMessage> 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<MimeMessage> 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<MimeMessage> filteredMessages = new LinkedList<MimeMessage>();
List<MimeMessage> 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

View File

@@ -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<Advice> 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<Advice> 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<Object>) 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(

View File

@@ -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();
}