INT-4147: Fix NPE in the FileReadingMessageSource
JIRA: https://jira.spring.io/browse/INT-4147 By default `FileReadingMessageSource` is created without `filter`, at the same time internal `WatchServiceDirectoryScanner` is supplied with default `filter` by its `DefaultDirectoryScanner` super class. A `DELETE` watch event condition around `FileReadingMessageSource.this.filter` is wrong, because it is `null` by default. Therefore `DELETE` events never succeed * Modify `FileInboundTransactionTests` to accept `DELETE` watch event as well and verify that `ResettableFileListFilter.remove(File)` is performed * Call newly created getter for `DefaultDirectoryScanner.filter` in the `WatchServiceDirectoryScanner` to assert against supplied `filter` Polishing **Cherry-pick to 4.3.x** Conflicts: spring-integration-file/src/main/java/org/springframework/integration/file/DefaultDirectoryScanner.java Resolved.
This commit is contained in:
committed by
Gary Russell
parent
80b1f4c7bc
commit
4653090686
@@ -21,11 +21,11 @@ import java.util.ArrayList;
|
||||
import java.util.Arrays;
|
||||
import java.util.List;
|
||||
|
||||
import org.springframework.messaging.MessagingException;
|
||||
import org.springframework.integration.file.filters.AcceptOnceFileListFilter;
|
||||
import org.springframework.integration.file.filters.CompositeFileListFilter;
|
||||
import org.springframework.integration.file.filters.FileListFilter;
|
||||
import org.springframework.integration.file.filters.IgnoreHiddenFileListFilter;
|
||||
import org.springframework.messaging.MessagingException;
|
||||
|
||||
/**
|
||||
* Default directory scanner and base class for other directory scanners.
|
||||
@@ -41,17 +41,6 @@ public class DefaultDirectoryScanner implements DirectoryScanner {
|
||||
|
||||
private volatile FileLocker locker;
|
||||
|
||||
public void setFilter(FileListFilter<File> filter) {
|
||||
this.filter = filter;
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritDoc}
|
||||
*/
|
||||
public final void setLocker(FileLocker locker) {
|
||||
this.locker = locker;
|
||||
}
|
||||
|
||||
/**
|
||||
* Initializes {@link DefaultDirectoryScanner#filter} with a default list of
|
||||
* {@link FileListFilter}s using a {@link CompositeFileListFilter}:
|
||||
@@ -67,16 +56,36 @@ public class DefaultDirectoryScanner implements DirectoryScanner {
|
||||
this.filter = new CompositeFileListFilter<File>(defaultFilters);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setFilter(FileListFilter<File> filter) {
|
||||
this.filter = filter;
|
||||
}
|
||||
|
||||
protected FileListFilter<File> getFilter() {
|
||||
return this.filter;
|
||||
}
|
||||
|
||||
@Override
|
||||
public final void setLocker(FileLocker locker) {
|
||||
this.locker = locker;
|
||||
}
|
||||
|
||||
protected FileLocker getLocker() {
|
||||
return this.locker;
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritDoc}
|
||||
* <p>
|
||||
* This class takes the minimal implementation and merely delegates to the
|
||||
* locker if set.
|
||||
*/
|
||||
@Override
|
||||
public final boolean tryClaim(File file) {
|
||||
return (this.locker == null) || this.locker.lock(file);
|
||||
}
|
||||
|
||||
@Override
|
||||
public final List<File> listFiles(File directory) throws IllegalArgumentException {
|
||||
File[] files = listEligibleFiles(directory);
|
||||
if (files == null) {
|
||||
|
||||
@@ -498,8 +498,8 @@ public class FileReadingMessageSource extends IntegrationObjectSupport implement
|
||||
}
|
||||
|
||||
if (event.kind() == StandardWatchEventKinds.ENTRY_DELETE) {
|
||||
if (FileReadingMessageSource.this.filter instanceof ResettableFileListFilter) {
|
||||
((ResettableFileListFilter<File>) FileReadingMessageSource.this.filter).remove(file);
|
||||
if (getFilter() instanceof ResettableFileListFilter) {
|
||||
((ResettableFileListFilter<File>) getFilter()).remove(file);
|
||||
}
|
||||
boolean fileRemoved = files.remove(file);
|
||||
if (fileRemoved && logger.isDebugEnabled()) {
|
||||
|
||||
@@ -14,7 +14,8 @@
|
||||
<int-file:inbound-channel-adapter id="pseudoTx"
|
||||
channel="input" auto-startup="false"
|
||||
directory="#{T (org.springframework.integration.file.FileInboundTransactionTests).tmpDir.root}/si-test1"
|
||||
use-watch-service="true">
|
||||
use-watch-service="true"
|
||||
watch-events="CREATE,DELETE">
|
||||
<int:poller fixed-rate="500">
|
||||
<int:transactional synchronization-factory="syncFactoryA"/>
|
||||
</int:poller>
|
||||
|
||||
@@ -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.MessageHandler;
|
||||
@@ -85,6 +89,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<File> 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(new MessageHandler() {
|
||||
@@ -115,8 +129,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
|
||||
|
||||
Reference in New Issue
Block a user