From 59105ef8d120646d085b332f90a2877c7d7e03b5 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Mon, 19 Nov 2012 18:36:54 -0500 Subject: [PATCH] INT-2821 Find Existing Mails if No RECENT Support Previously, if an IMAP server supports IDLE, but not RECENT, no existing messages were retrieved until a new message arrived. INT-2821 Polishing - Fix Race Condition Use searchForMessages() when server doesn't support RECENT. Simply skipping the first idle() doesn't work because receive() closes the folder, there was a race condition where a new message could arrive between the previous receive() and folder.open() in waitForMessages(). Updated mock test so that first call to waitForMessages() finds messages and does not idle(), and subsequent calls does not find messages and goes to idle(). Also tested in the debugger against gmail which (at the time of writing) does not support RECENT. --- .../integration/mail/ImapMailReceiver.java | 12 +- .../mail/ImapMailReceiverTests.java | 262 ++++++++++++++---- 2 files changed, 220 insertions(+), 54 deletions(-) diff --git a/spring-integration-mail/src/main/java/org/springframework/integration/mail/ImapMailReceiver.java b/spring-integration-mail/src/main/java/org/springframework/integration/mail/ImapMailReceiver.java index 05b9799c42..f97eef581b 100755 --- a/spring-integration-mail/src/main/java/org/springframework/integration/mail/ImapMailReceiver.java +++ b/spring-integration-mail/src/main/java/org/springframework/integration/mail/ImapMailReceiver.java @@ -41,10 +41,11 @@ import com.sun.mail.imap.IMAPMessage; * the option of blocking until new messages are available prior to calling * {@link #receive()}. That option is only available if the server supports * the {@link IMAPFolder#idle() idle} command. - * + * * @author Arjen Poutsma * @author Mark Fisher * @author Oleg Zhurakousky + * @author Gary Russell */ public class ImapMailReceiver extends AbstractMailReceiver { @@ -95,6 +96,11 @@ public class ImapMailReceiver extends AbstractMailReceiver { if (imapFolder.hasNewMessages()) { return; } + else if (!imapFolder.getPermanentFlags().contains(Flags.Flag.RECENT)) { + if (searchForNewMessages().length > 0) { + return; + } + } imapFolder.addMessageCountListener(this.messageCountListener); try { imapFolder.idle(); @@ -174,7 +180,7 @@ public class ImapMailReceiver extends AbstractMailReceiver { "USER flags which will be used to prevent duplicates during email fetch."); Flags siFlags = new Flags(); siFlags.add(SI_USER_FLAG); - notFlagged = new NotTerm(new FlagTerm(siFlags, true)); + notFlagged = new NotTerm(new FlagTerm(siFlags, true)); } else { logger.debug("This email server does not support RECENT or USER flags. " + @@ -191,6 +197,7 @@ public class ImapMailReceiver extends AbstractMailReceiver { return searchTerm; } + @Override protected void setAdditionalFlags(Message message) throws MessagingException { super.setAdditionalFlags(message); if (this.shouldMarkMessagesAsRead) { @@ -204,6 +211,7 @@ public class ImapMailReceiver extends AbstractMailReceiver { */ private static class SimpleMessageCountListener extends MessageCountAdapter { + @Override public void messagesAdded(MessageCountEvent event) { Message[] messages = event.getMessages(); for (Message message : messages) { 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 962518e4b0..4196c19c63 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 @@ -17,6 +17,8 @@ package org.springframework.integration.mail; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertNull; +import static org.junit.Assert.assertTrue; import static org.mockito.Mockito.doAnswer; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.spy; @@ -26,6 +28,8 @@ import static org.mockito.Mockito.when; import java.lang.reflect.Field; import java.util.Properties; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; import javax.mail.Flags; @@ -42,7 +46,6 @@ import org.junit.Test; import org.mockito.Mockito; import org.mockito.invocation.InvocationOnMock; import org.mockito.stubbing.Answer; - import org.springframework.beans.DirectFieldAccessor; import org.springframework.context.ApplicationContext; import org.springframework.context.support.ClassPathXmlApplicationContext; @@ -62,9 +65,9 @@ import com.sun.mail.imap.IMAPFolder; * */ public class ImapMailReceiverTests { - + private AtomicInteger failed = new AtomicInteger(0); - + @Test public void receiveAndMarkAsReadDontDelete() throws Exception{ AbstractMailReceiver receiver = new ImapMailReceiver(); @@ -80,7 +83,7 @@ public class ImapMailReceiverTests { Message msg1 = mock(MimeMessage.class); Message msg2 = mock(MimeMessage.class); final Message[] messages = new Message[]{msg1, msg2}; - + doAnswer(new Answer() { public Object answer(InvocationOnMock invocation) throws Throwable { DirectFieldAccessor accessor = new DirectFieldAccessor(invocation.getMock()); @@ -88,17 +91,17 @@ public class ImapMailReceiverTests { if (folderOpenMode != Folder.READ_WRITE){ throw new IllegalArgumentException("Folder had to be open in READ_WRITE mode"); } - + return null; } }).when(receiver).openFolder(); - + doAnswer(new Answer() { public Object answer(InvocationOnMock invocation) throws Throwable { return messages; } }).when(receiver).searchForNewMessages(); - + doAnswer(new Answer() { public Object answer(InvocationOnMock invocation) throws Throwable { return null; @@ -116,13 +119,13 @@ public class ImapMailReceiverTests { receiver.setShouldDeleteMessages(true); receiver = spy(receiver); receiver.afterPropertiesSet(); - + Field folderField = AbstractMailReceiver.class.getDeclaredField("folder"); folderField.setAccessible(true); Folder folder = mock(Folder.class); when(folder.getPermanentFlags()).thenReturn(new Flags(Flags.Flag.USER)); folderField.set(receiver, folder); - + Message msg1 = mock(MimeMessage.class); Message msg2 = mock(MimeMessage.class); final Message[] messages = new Message[]{msg1, msg2}; @@ -136,13 +139,13 @@ public class ImapMailReceiverTests { return null; } }).when(receiver).openFolder(); - + doAnswer(new Answer() { public Object answer(InvocationOnMock invocation) throws Throwable { return messages; } }).when(receiver).searchForNewMessages(); - + doAnswer(new Answer() { public Object answer(InvocationOnMock invocation) throws Throwable { return null; @@ -159,14 +162,14 @@ public class ImapMailReceiverTests { ((ImapMailReceiver)receiver).setShouldMarkMessagesAsRead(false); receiver = spy(receiver); receiver.afterPropertiesSet(); - + Field folderField = AbstractMailReceiver.class.getDeclaredField("folder"); folderField.setAccessible(true); Folder folder = mock(Folder.class); when(folder.getPermanentFlags()).thenReturn(new Flags(Flags.Flag.USER)); folderField.set(receiver, folder); - - + + Message msg1 = mock(MimeMessage.class); Message msg2 = mock(MimeMessage.class); final Message[] messages = new Message[]{msg1, msg2}; @@ -175,13 +178,13 @@ public class ImapMailReceiverTests { return null; } }).when(receiver).openFolder(); - + doAnswer(new Answer() { public Object answer(InvocationOnMock invocation) throws Throwable { return messages; } }).when(receiver).searchForNewMessages(); - + doAnswer(new Answer() { public Object answer(InvocationOnMock invocation) throws Throwable { return null; @@ -199,13 +202,13 @@ public class ImapMailReceiverTests { ((ImapMailReceiver)receiver).setShouldMarkMessagesAsRead(false); receiver = spy(receiver); receiver.afterPropertiesSet(); - + Field folderField = AbstractMailReceiver.class.getDeclaredField("folder"); folderField.setAccessible(true); Folder folder = mock(Folder.class); when(folder.getPermanentFlags()).thenReturn(new Flags(Flags.Flag.USER)); folderField.set(receiver, folder); - + Message msg1 = mock(MimeMessage.class); Message msg2 = mock(MimeMessage.class); final Message[] messages = new Message[]{msg1, msg2}; @@ -219,13 +222,13 @@ public class ImapMailReceiverTests { return null; } }).when(receiver).openFolder(); - + doAnswer(new Answer() { public Object answer(InvocationOnMock invocation) throws Throwable { return messages; } }).when(receiver).searchForNewMessages(); - + doAnswer(new Answer() { public Object answer(InvocationOnMock invocation) throws Throwable { return null; @@ -243,13 +246,13 @@ public class ImapMailReceiverTests { AbstractMailReceiver receiver = new ImapMailReceiver(); receiver = spy(receiver); receiver.afterPropertiesSet(); - + Field folderField = AbstractMailReceiver.class.getDeclaredField("folder"); folderField.setAccessible(true); Folder folder = mock(Folder.class); when(folder.getPermanentFlags()).thenReturn(new Flags(Flags.Flag.USER)); folderField.set(receiver, folder); - + Message msg1 = mock(MimeMessage.class); Message msg2 = mock(MimeMessage.class); final Message[] messages = new Message[]{msg1, msg2}; @@ -263,13 +266,13 @@ public class ImapMailReceiverTests { return null; } }).when(receiver).openFolder(); - + doAnswer(new Answer() { public Object answer(InvocationOnMock invocation) throws Throwable { return messages; } }).when(receiver).searchForNewMessages(); - + doAnswer(new Answer() { public Object answer(InvocationOnMock invocation) throws Throwable { return null; @@ -283,22 +286,22 @@ public class ImapMailReceiverTests { @Test @Ignore public void testMessageHistory() throws Exception{ - ApplicationContext context = + ApplicationContext context = new ClassPathXmlApplicationContext("ImapIdleChannelAdapterParserTests-context.xml", ImapIdleChannelAdapterParserTests.class); ImapIdleChannelAdapter adapter = context.getBean("simpleAdapter", ImapIdleChannelAdapter.class); - + AbstractMailReceiver receiver = new ImapMailReceiver(); receiver = spy(receiver); receiver.afterPropertiesSet(); - + DirectFieldAccessor adapterAccessor = new DirectFieldAccessor(adapter); adapterAccessor.setPropertyValue("mailReceiver", receiver); - + MimeMessage mailMessage = mock(MimeMessage.class); Flags flags = mock(Flags.class); when(mailMessage.getFlags()).thenReturn(flags); final Message[] messages = new Message[]{mailMessage}; - + doAnswer(new Answer() { public Object answer(InvocationOnMock invocation) throws Throwable { DirectFieldAccessor accesor = new DirectFieldAccessor((invocation.getMock())); @@ -308,19 +311,19 @@ public class ImapMailReceiverTests { return null; } }).when(receiver).openFolder(); - + doAnswer(new Answer() { public Object answer(InvocationOnMock invocation) throws Throwable { return messages; } }).when(receiver).searchForNewMessages(); - + doAnswer(new Answer() { public Object answer(InvocationOnMock invocation) throws Throwable { return null; } }).when(receiver).fetchMessages(messages); - + PollableChannel channel = context.getBean("channel", PollableChannel.class); adapter.start(); @@ -334,16 +337,17 @@ public class ImapMailReceiverTests { @Test public void testIdleChannelAdapterException() throws Exception{ - ApplicationContext context = + ApplicationContext context = new ClassPathXmlApplicationContext("ImapIdleChannelAdapterParserTests-context.xml", ImapIdleChannelAdapterParserTests.class); ImapIdleChannelAdapter adapter = context.getBean("simpleAdapter", ImapIdleChannelAdapter.class); //ImapMailReceiver receiver = (ImapMailReceiver) TestUtils.getPropertyValue(adapter, "mailReceiver"); - - - + + + DirectChannel channel = new DirectChannel(); channel.subscribe(new AbstractReplyProducingMessageHandler() { + @Override protected Object handleRequestMessage(org.springframework.integration.Message requestMessage) { throw new RuntimeException("Failed"); } @@ -351,58 +355,212 @@ public class ImapMailReceiverTests { adapter.setOutputChannel(channel); QueueChannel errorChannel = new QueueChannel(); adapter.setErrorChannel(errorChannel); - + AbstractMailReceiver receiver = new ImapMailReceiver(); receiver = spy(receiver); receiver.afterPropertiesSet(); - + Field folderField = AbstractMailReceiver.class.getDeclaredField("folder"); folderField.setAccessible(true); Folder folder = mock(IMAPFolder.class); when(folder.getPermanentFlags()).thenReturn(new Flags(Flags.Flag.USER)); folderField.set(receiver, folder); - + doAnswer(new Answer() { public Object answer(InvocationOnMock invocation) throws Throwable { return true; } }).when(folder).isOpen(); - + doAnswer(new Answer() { public Object answer(InvocationOnMock invocation) throws Throwable { return null; } }).when(receiver).openFolder(); - + DirectFieldAccessor adapterAccessor = new DirectFieldAccessor(adapter); adapterAccessor.setPropertyValue("mailReceiver", receiver); - + MimeMessage mailMessage = mock(MimeMessage.class); Flags flags = mock(Flags.class); when(mailMessage.getFlags()).thenReturn(flags); final Message[] messages = new Message[]{mailMessage}; - + doAnswer(new Answer() { public Object answer(InvocationOnMock invocation) throws Throwable { return messages; } }).when(receiver).searchForNewMessages(); - + doAnswer(new Answer() { public Object answer(InvocationOnMock invocation) throws Throwable { return null; } }).when(receiver).fetchMessages(messages); - + adapter.start(); org.springframework.integration.Message replMessage = errorChannel.receive(10000); assertNotNull(replMessage); assertEquals("Failed", ((Exception) replMessage.getPayload()).getCause().getMessage()); } - + + @Test + public void testNoInitialIdleDelayWhenRecentNotSupported() throws Exception{ + ApplicationContext context = + new ClassPathXmlApplicationContext("ImapIdleChannelAdapterParserTests-context.xml", ImapIdleChannelAdapterParserTests.class); + ImapIdleChannelAdapter adapter = context.getBean("simpleAdapter", ImapIdleChannelAdapter.class); + + QueueChannel channel = new QueueChannel(); + adapter.setOutputChannel(channel); + + ImapMailReceiver receiver = new ImapMailReceiver("imap:foo"); + receiver = spy(receiver); + receiver.afterPropertiesSet(); + + final IMAPFolder folder = mock(IMAPFolder.class); + when(folder.getPermanentFlags()).thenReturn(new Flags(Flags.Flag.USER)); + when(folder.isOpen()).thenReturn(false).thenReturn(true); + when(folder.exists()).thenReturn(true); + + DirectFieldAccessor adapterAccessor = new DirectFieldAccessor(adapter); + adapterAccessor.setPropertyValue("mailReceiver", receiver); + + Field storeField = AbstractMailReceiver.class.getDeclaredField("store"); + storeField.setAccessible(true); + Store store = mock(Store.class); + when(store.isConnected()).thenReturn(true); + when(store.getFolder(Mockito.any(URLName.class))).thenReturn(folder); + storeField.set(receiver, store); + + doAnswer(new Answer() { + public Object answer(InvocationOnMock invocation) throws Throwable { + return folder; + } + }).when(receiver).getFolder(); + + MimeMessage mailMessage = mock(MimeMessage.class); + Flags flags = mock(Flags.class); + when(mailMessage.getFlags()).thenReturn(flags); + final Message[] messages = new Message[]{mailMessage}; + + final AtomicInteger shouldFindMessagesCounter = new AtomicInteger(2); + doAnswer(new Answer() { + public Object answer(InvocationOnMock invocation) throws Throwable { + /* + * Return the message from first invocation of waitForMessages() + * and in receive(); then return false in the next call to + * waitForMessages() so we enter idle(); counter will be reset + * to 1 in the mocked idle(). + */ + if (shouldFindMessagesCounter.decrementAndGet() >= 0) { + return messages; + } + else { + return new Message[0]; + } + } + }).when(receiver).searchForNewMessages(); + + doAnswer(new Answer() { + public Object answer(InvocationOnMock invocation) throws Throwable { + return null; + } + }).when(receiver).fetchMessages(messages); + + doAnswer(new Answer() { + public Object answer(InvocationOnMock invocation) throws Throwable { + Thread.sleep(5000); + shouldFindMessagesCounter.set(1); + return null; + } + }).when(folder).idle(); + + adapter.start(); + + /* + * Idle takes 5 seconds; if all is well, we should receive the first message + * before then. + */ + assertNotNull(channel.receive(3000)); + // We should not receive any more until the next idle elapses + assertNull(channel.receive(3000)); + assertNotNull(channel.receive(6000)); + } + + @Test + public void testInitialIdleDelayWhenRecentIsSupported() throws Exception{ + ApplicationContext context = + new ClassPathXmlApplicationContext("ImapIdleChannelAdapterParserTests-context.xml", ImapIdleChannelAdapterParserTests.class); + ImapIdleChannelAdapter adapter = context.getBean("simpleAdapter", ImapIdleChannelAdapter.class); + + QueueChannel channel = new QueueChannel(); + adapter.setOutputChannel(channel); + + ImapMailReceiver receiver = new ImapMailReceiver("imap:foo"); + receiver = spy(receiver); + receiver.afterPropertiesSet(); + + final IMAPFolder folder = mock(IMAPFolder.class); + when(folder.getPermanentFlags()).thenReturn(new Flags(Flags.Flag.RECENT)); + when(folder.isOpen()).thenReturn(false).thenReturn(true); + when(folder.exists()).thenReturn(true); + + DirectFieldAccessor adapterAccessor = new DirectFieldAccessor(adapter); + adapterAccessor.setPropertyValue("mailReceiver", receiver); + + Field storeField = AbstractMailReceiver.class.getDeclaredField("store"); + storeField.setAccessible(true); + Store store = mock(Store.class); + when(store.isConnected()).thenReturn(true); + when(store.getFolder(Mockito.any(URLName.class))).thenReturn(folder); + storeField.set(receiver, store); + + doAnswer(new Answer() { + public Object answer(InvocationOnMock invocation) throws Throwable { + return folder; + } + }).when(receiver).getFolder(); + + MimeMessage mailMessage = mock(MimeMessage.class); + Flags flags = mock(Flags.class); + when(mailMessage.getFlags()).thenReturn(flags); + final Message[] messages = new Message[]{mailMessage}; + + doAnswer(new Answer() { + public Object answer(InvocationOnMock invocation) throws Throwable { + return messages; + } + }).when(receiver).searchForNewMessages(); + + doAnswer(new Answer() { + public Object answer(InvocationOnMock invocation) throws Throwable { + return null; + } + }).when(receiver).fetchMessages(messages); + + final CountDownLatch idles = new CountDownLatch(2); + doAnswer(new Answer() { + public Object answer(InvocationOnMock invocation) throws Throwable { + idles.countDown(); + Thread.sleep(5000); + return null; + } + }).when(folder).idle(); + + adapter.start(); + + /* + * Idle takes 5 seconds; since this server supports RECENT, we should + * not receive any early messages. + */ + assertNull(channel.receive(3000)); + assertNotNull(channel.receive(5000)); + assertTrue(idles.await(5, TimeUnit.SECONDS)); + } + @Test // see INT-1801 public void testImapLifecycleForRaceCondition() throws Exception{ - + for (int i = 0; i < 1000; i++) { final ImapMailReceiver receiver = new ImapMailReceiver("imap://foo"); Store store = mock(Store.class); @@ -412,12 +570,12 @@ public class ImapMailReceiverTests { when(folder.search((SearchTerm) Mockito.any())).thenReturn(new Message[]{}); when(store.getFolder(Mockito.any(URLName.class))).thenReturn(folder); when(folder.getPermanentFlags()).thenReturn(new Flags(Flags.Flag.USER)); - - + + DirectFieldAccessor df = new DirectFieldAccessor(receiver); df.setPropertyValue("store", store); receiver.afterPropertiesSet(); - + new Thread(new Runnable() { public void run(){ try { @@ -428,10 +586,10 @@ public class ImapMailReceiverTests { failed.getAndIncrement(); } } - + } }).start(); - + new Thread(new Runnable() { public void run(){ try {