diff --git a/spring-integration-file/src/main/java/org/springframework/integration/file/DefaultDirectoryScanner.java b/spring-integration-file/src/main/java/org/springframework/integration/file/DefaultDirectoryScanner.java index f998c0e6e6..c75ef55c2e 100644 --- a/spring-integration-file/src/main/java/org/springframework/integration/file/DefaultDirectoryScanner.java +++ b/spring-integration-file/src/main/java/org/springframework/integration/file/DefaultDirectoryScanner.java @@ -62,11 +62,19 @@ public class DefaultDirectoryScanner implements DirectoryScanner { this.filter = filter; } + protected FileListFilter getFilter() { + return this.filter; + } + @Override public final void setLocker(FileLocker locker) { this.locker = locker; } + protected FileLocker getLocker() { + return this.locker; + } + /** * This class takes the minimal implementation and merely delegates to the locker if set. * @param file the file to try to claim. diff --git a/spring-integration-file/src/main/java/org/springframework/integration/file/FileReadingMessageSource.java b/spring-integration-file/src/main/java/org/springframework/integration/file/FileReadingMessageSource.java index fda2b06aad..2463832d1f 100644 --- a/spring-integration-file/src/main/java/org/springframework/integration/file/FileReadingMessageSource.java +++ b/spring-integration-file/src/main/java/org/springframework/integration/file/FileReadingMessageSource.java @@ -490,8 +490,8 @@ public class FileReadingMessageSource extends IntegrationObjectSupport implement } if (event.kind() == StandardWatchEventKinds.ENTRY_DELETE) { - if (FileReadingMessageSource.this.filter instanceof ResettableFileListFilter) { - ((ResettableFileListFilter) FileReadingMessageSource.this.filter).remove(file); + if (getFilter() != null && getFilter() instanceof ResettableFileListFilter) { + ((ResettableFileListFilter) getFilter()).remove(file); } boolean fileRemoved = files.remove(file); if (fileRemoved && logger.isDebugEnabled()) { diff --git a/spring-integration-file/src/test/java/org/springframework/integration/file/FileInboundTransactionTests-context.xml b/spring-integration-file/src/test/java/org/springframework/integration/file/FileInboundTransactionTests-context.xml index ad143397a1..1b9e01d22f 100644 --- a/spring-integration-file/src/test/java/org/springframework/integration/file/FileInboundTransactionTests-context.xml +++ b/spring-integration-file/src/test/java/org/springframework/integration/file/FileInboundTransactionTests-context.xml @@ -14,7 +14,8 @@ + use-watch-service="true" + watch-events="CREATE,DELETE"> diff --git a/spring-integration-file/src/test/java/org/springframework/integration/file/FileInboundTransactionTests.java b/spring-integration-file/src/test/java/org/springframework/integration/file/FileInboundTransactionTests.java index 16a3ac64fc..a0326fadaa 100644 --- a/spring-integration-file/src/test/java/org/springframework/integration/file/FileInboundTransactionTests.java +++ b/spring-integration-file/src/test/java/org/springframework/integration/file/FileInboundTransactionTests.java @@ -22,6 +22,8 @@ import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertNotNull; import static org.junit.Assert.assertThat; import static org.junit.Assert.assertTrue; +import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.verify; import java.io.File; import java.util.concurrent.CountDownLatch; @@ -32,8 +34,10 @@ import org.junit.Test; import org.junit.rules.TemporaryFolder; import org.junit.runner.RunWith; +import org.springframework.beans.DirectFieldAccessor; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.integration.endpoint.SourcePollingChannelAdapter; +import org.springframework.integration.file.filters.ResettableFileListFilter; import org.springframework.integration.test.util.TestUtils; import org.springframework.messaging.Message; import org.springframework.messaging.MessagingException; @@ -84,6 +88,16 @@ public class FileInboundTransactionTests { @Test public void testNoTx() throws Exception { + + Object scanner = TestUtils.getPropertyValue(pseudoTx.getMessageSource(), "scanner"); + assertThat(scanner.getClass().getName(), containsString("FileReadingMessageSource$WatchServiceDirectoryScanner")); + + @SuppressWarnings("unchecked") + ResettableFileListFilter fileListFilter = + spy(TestUtils.getPropertyValue(scanner, "filter", ResettableFileListFilter.class)); + + new DirectFieldAccessor(scanner).setPropertyValue("filter", fileListFilter); + final CountDownLatch latch = new CountDownLatch(1); final AtomicBoolean crash = new AtomicBoolean(); input.subscribe(message -> { @@ -110,8 +124,7 @@ public class FileInboundTransactionTests { assertFalse(transactionManager.getCommitted()); assertFalse(transactionManager.getRolledBack()); - Object scanner = TestUtils.getPropertyValue(pseudoTx.getMessageSource(), "scanner"); - assertThat(scanner.getClass().getName(), containsString("FileReadingMessageSource$WatchServiceDirectoryScanner")); + verify(fileListFilter).remove(new File(tmpDir.getRoot(), "si-test1/foo")); } @Test