Close file resource after (S)FTP streaming tests

Windows is too scrupulous for non closed file resources.
Our `RemoteFileTestSupport` recreates a test directory with files
for each test.
When we don't close some file handler, we are not able to delete the
directory and we fail with subsequent tests.

* Close `InputStream` for each polled file in new tests in the
`SftpStreamingMessageSourceTests`

* Close `CLOSEABLE_RESOURCE` in new tests in the
`FtpStreamingMessageSourceTests`
* Rework `FtpStreamingMessageSourceTests.testAllContents()` do not
poll all the files on each polling cycle - this causes a race condition
when we don't close `CLOSEABLE_RESOURCE` yet in the `StreamTransformer`,
but try to proceed with recreation a test directory structure

**Cherry-pick to 5.0.x**
This commit is contained in:
Artem Bilan
2018-06-29 21:28:28 -04:00
parent 6affb3bb2f
commit 6ecd948a32
2 changed files with 48 additions and 17 deletions

View File

@@ -17,16 +17,19 @@
package org.springframework.integration.ftp.inbound; package org.springframework.integration.ftp.inbound;
import static org.hamcrest.CoreMatchers.containsString; import static org.hamcrest.CoreMatchers.containsString;
import static org.hamcrest.CoreMatchers.instanceOf;
import static org.hamcrest.Matchers.equalTo; import static org.hamcrest.Matchers.equalTo;
import static org.hamcrest.Matchers.instanceOf;
import static org.junit.Assert.assertNotNull; import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.assertThat; import static org.junit.Assert.assertThat;
import java.io.Closeable;
import java.io.IOException;
import java.io.InputStream; import java.io.InputStream;
import java.util.Comparator; import java.util.Comparator;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import org.apache.commons.net.ftp.FTPFile; import org.apache.commons.net.ftp.FTPFile;
import org.junit.Rule;
import org.junit.Test; import org.junit.Test;
import org.junit.runner.RunWith; import org.junit.runner.RunWith;
@@ -34,6 +37,7 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.ApplicationContext; import org.springframework.context.ApplicationContext;
import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Configuration;
import org.springframework.integration.StaticMessageHeaderAccessor;
import org.springframework.integration.annotation.InboundChannelAdapter; import org.springframework.integration.annotation.InboundChannelAdapter;
import org.springframework.integration.annotation.Transformer; import org.springframework.integration.annotation.Transformer;
import org.springframework.integration.channel.QueueChannel; import org.springframework.integration.channel.QueueChannel;
@@ -45,10 +49,11 @@ import org.springframework.integration.file.filters.AcceptAllFileListFilter;
import org.springframework.integration.file.remote.FileInfo; import org.springframework.integration.file.remote.FileInfo;
import org.springframework.integration.file.remote.session.SessionFactory; import org.springframework.integration.file.remote.session.SessionFactory;
import org.springframework.integration.ftp.FtpTestSupport; import org.springframework.integration.ftp.FtpTestSupport;
import org.springframework.integration.ftp.filters.FtpPersistentAcceptOnceFileListFilter;
import org.springframework.integration.ftp.session.FtpFileInfo; import org.springframework.integration.ftp.session.FtpFileInfo;
import org.springframework.integration.ftp.session.FtpRemoteFileTemplate; import org.springframework.integration.ftp.session.FtpRemoteFileTemplate;
import org.springframework.integration.metadata.SimpleMetadataStore;
import org.springframework.integration.scheduling.PollerMetadata; import org.springframework.integration.scheduling.PollerMetadata;
import org.springframework.integration.test.rule.Log4j2LevelAdjuster;
import org.springframework.integration.transformer.StreamTransformer; import org.springframework.integration.transformer.StreamTransformer;
import org.springframework.messaging.Message; import org.springframework.messaging.Message;
import org.springframework.scheduling.support.PeriodicTrigger; import org.springframework.scheduling.support.PeriodicTrigger;
@@ -81,10 +86,8 @@ public class FtpStreamingMessageSourceTests extends FtpTestSupport {
@Autowired @Autowired
private ApplicationContext context; private ApplicationContext context;
@Rule @Autowired
public Log4j2LevelAdjuster adjuster = private ConcurrentMap<String, String> metadataMap;
Log4j2LevelAdjuster.debug()
.categories(true, "org.apache.commons");
@SuppressWarnings("unchecked") @SuppressWarnings("unchecked")
@Test @Test
@@ -116,34 +119,46 @@ public class FtpStreamingMessageSourceTests extends FtpTestSupport {
this.adapter.stop(); this.adapter.stop();
this.source.setFileInfoJson(false); this.source.setFileInfoJson(false);
this.data.purge(null); this.data.purge(null);
this.metadataMap.clear();
this.adapter.start(); this.adapter.start();
received = (Message<byte[]>) this.data.receive(10000); received = (Message<byte[]>) this.data.receive(10000);
assertNotNull(received); assertNotNull(received);
assertThat(received.getHeaders().get(FileHeaders.REMOTE_FILE_INFO), instanceOf(FtpFileInfo.class));
this.adapter.stop(); this.adapter.stop();
assertThat(received.getHeaders().get(FileHeaders.REMOTE_FILE_INFO), instanceOf(FtpFileInfo.class));
} }
@Test @Test
public void testMaxFetch() { public void testMaxFetch() throws IOException {
FtpStreamingMessageSource messageSource = buildsource(); FtpStreamingMessageSource messageSource = buildSource();
messageSource.setFilter(new AcceptAllFileListFilter<>()); messageSource.setFilter(new AcceptAllFileListFilter<>());
messageSource.afterPropertiesSet(); messageSource.afterPropertiesSet();
Message<InputStream> received = messageSource.receive(); Message<InputStream> received = messageSource.receive();
assertNotNull(received); assertNotNull(received);
assertThat(received.getHeaders().get(FileHeaders.REMOTE_FILE), equalTo(" ftpSource1.txt")); assertThat(received.getHeaders().get(FileHeaders.REMOTE_FILE), equalTo(" ftpSource1.txt"));
Closeable closeableResource = StaticMessageHeaderAccessor.getCloseableResource(received);
if (closeableResource != null) {
closeableResource.close();
}
} }
@Test @Test
public void testMaxFetchNoFilter() { public void testMaxFetchNoFilter() throws IOException {
FtpStreamingMessageSource messageSource = buildsource(); FtpStreamingMessageSource messageSource = buildSource();
messageSource.setFilter(null); messageSource.setFilter(null);
messageSource.afterPropertiesSet(); messageSource.afterPropertiesSet();
Message<InputStream> received = messageSource.receive(); Message<InputStream> received = messageSource.receive();
assertNotNull(received); assertNotNull(received);
assertThat(received.getHeaders().get(FileHeaders.REMOTE_FILE), equalTo(" ftpSource1.txt")); assertThat(received.getHeaders().get(FileHeaders.REMOTE_FILE), equalTo(" ftpSource1.txt"));
Closeable closeableResource = StaticMessageHeaderAccessor.getCloseableResource(received);
if (closeableResource != null) {
closeableResource.close();
}
} }
private FtpStreamingMessageSource buildsource() { private FtpStreamingMessageSource buildSource() {
FtpStreamingMessageSource messageSource = new FtpStreamingMessageSource(this.config.template(), FtpStreamingMessageSource messageSource = new FtpStreamingMessageSource(this.config.template(),
Comparator.comparing(FileInfo::getFilename)); Comparator.comparing(FileInfo::getFilename));
messageSource.setRemoteDirectory("ftpSource/"); messageSource.setRemoteDirectory("ftpSource/");
@@ -169,12 +184,20 @@ public class FtpStreamingMessageSourceTests extends FtpTestSupport {
return pollerMetadata; return pollerMetadata;
} }
@Bean
public ConcurrentMap<String, String> metadataMap() {
return new ConcurrentHashMap<>();
}
@Bean @Bean
@InboundChannelAdapter(channel = "stream", autoStartup = "false") @InboundChannelAdapter(channel = "stream", autoStartup = "false")
public MessageSource<InputStream> ftpMessageSource() { public MessageSource<InputStream> ftpMessageSource() {
FtpStreamingMessageSource messageSource = new FtpStreamingMessageSource(template(), FtpStreamingMessageSource messageSource = new FtpStreamingMessageSource(template(),
Comparator.comparing(FileInfo::getFilename)); Comparator.comparing(FileInfo::getFilename));
messageSource.setFilter(new AcceptAllFileListFilter<>()); messageSource.setFilter(
new FtpPersistentAcceptOnceFileListFilter(
new SimpleMetadataStore(metadataMap()), "testStreaming"));
messageSource.setRemoteDirectory("ftpSource/"); messageSource.setRemoteDirectory("ftpSource/");
return messageSource; return messageSource;
} }

View File

@@ -23,6 +23,7 @@ import static org.hamcrest.Matchers.equalTo;
import static org.junit.Assert.assertNotNull; import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.assertThat; import static org.junit.Assert.assertThat;
import java.io.IOException;
import java.io.InputStream; import java.io.InputStream;
import java.util.Arrays; import java.util.Arrays;
import java.util.Comparator; import java.util.Comparator;
@@ -59,6 +60,7 @@ import com.jcraft.jsch.ChannelSftp.LsEntry;
/** /**
* @author Gary Russell * @author Gary Russell
* @author Artem Bilan * @author Artem Bilan
*
* @since 4.3 * @since 4.3
* *
*/ */
@@ -119,7 +121,7 @@ public class SftpStreamingMessageSourceTests extends SftpTestSupport {
} }
@Test @Test
public void testMaxFetch() { public void testMaxFetch() throws IOException {
SftpStreamingMessageSource messageSource = buildSource(); SftpStreamingMessageSource messageSource = buildSource();
messageSource.setFilter(new AcceptAllFileListFilter<>()); messageSource.setFilter(new AcceptAllFileListFilter<>());
messageSource.afterPropertiesSet(); messageSource.afterPropertiesSet();
@@ -127,10 +129,12 @@ public class SftpStreamingMessageSourceTests extends SftpTestSupport {
assertNotNull(received); assertNotNull(received);
assertThat(received.getHeaders().get(FileHeaders.REMOTE_FILE), assertThat(received.getHeaders().get(FileHeaders.REMOTE_FILE),
anyOf(equalTo(" sftpSource1.txt"), equalTo("sftpSource2.txt"))); anyOf(equalTo(" sftpSource1.txt"), equalTo("sftpSource2.txt")));
received.getPayload().close();
} }
@Test @Test
public void testMaxFetchNoFilter() { public void testMaxFetchNoFilter() throws IOException {
SftpStreamingMessageSource messageSource = buildSource(); SftpStreamingMessageSource messageSource = buildSource();
messageSource.setFilter(null); messageSource.setFilter(null);
messageSource.afterPropertiesSet(); messageSource.afterPropertiesSet();
@@ -138,10 +142,12 @@ public class SftpStreamingMessageSourceTests extends SftpTestSupport {
assertNotNull(received); assertNotNull(received);
assertThat(received.getHeaders().get(FileHeaders.REMOTE_FILE), assertThat(received.getHeaders().get(FileHeaders.REMOTE_FILE),
anyOf(equalTo(" sftpSource1.txt"), equalTo("sftpSource2.txt"))); anyOf(equalTo(" sftpSource1.txt"), equalTo("sftpSource2.txt")));
received.getPayload().close();
} }
@Test @Test
public void testMaxFetchLambdaFilter() { public void testMaxFetchLambdaFilter() throws IOException {
SftpStreamingMessageSource messageSource = buildSource(); SftpStreamingMessageSource messageSource = buildSource();
messageSource.setFilter(f -> Arrays.asList(f)); messageSource.setFilter(f -> Arrays.asList(f));
messageSource.afterPropertiesSet(); messageSource.afterPropertiesSet();
@@ -149,6 +155,8 @@ public class SftpStreamingMessageSourceTests extends SftpTestSupport {
assertNotNull(received); assertNotNull(received);
assertThat(received.getHeaders().get(FileHeaders.REMOTE_FILE), assertThat(received.getHeaders().get(FileHeaders.REMOTE_FILE),
anyOf(equalTo(" sftpSource1.txt"), equalTo("sftpSource2.txt"))); anyOf(equalTo(" sftpSource1.txt"), equalTo("sftpSource2.txt")));
received.getPayload().close();
} }
private SftpStreamingMessageSource buildSource() { private SftpStreamingMessageSource buildSource() {