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 {