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;
import static org.hamcrest.CoreMatchers.containsString;
import static org.hamcrest.CoreMatchers.instanceOf;
import static org.hamcrest.Matchers.equalTo;
import static org.hamcrest.Matchers.instanceOf;
import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.assertThat;
import java.io.Closeable;
import java.io.IOException;
import java.io.InputStream;
import java.util.Comparator;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import org.apache.commons.net.ftp.FTPFile;
import org.junit.Rule;
import org.junit.Test;
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.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.integration.StaticMessageHeaderAccessor;
import org.springframework.integration.annotation.InboundChannelAdapter;
import org.springframework.integration.annotation.Transformer;
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.session.SessionFactory;
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.FtpRemoteFileTemplate;
import org.springframework.integration.metadata.SimpleMetadataStore;
import org.springframework.integration.scheduling.PollerMetadata;
import org.springframework.integration.test.rule.Log4j2LevelAdjuster;
import org.springframework.integration.transformer.StreamTransformer;
import org.springframework.messaging.Message;
import org.springframework.scheduling.support.PeriodicTrigger;
@@ -81,10 +86,8 @@ public class FtpStreamingMessageSourceTests extends FtpTestSupport {
@Autowired
private ApplicationContext context;
@Rule
public Log4j2LevelAdjuster adjuster =
Log4j2LevelAdjuster.debug()
.categories(true, "org.apache.commons");
@Autowired
private ConcurrentMap<String, String> metadataMap;
@SuppressWarnings("unchecked")
@Test
@@ -116,34 +119,46 @@ public class FtpStreamingMessageSourceTests extends FtpTestSupport {
this.adapter.stop();
this.source.setFileInfoJson(false);
this.data.purge(null);
this.metadataMap.clear();
this.adapter.start();
received = (Message<byte[]>) this.data.receive(10000);
assertNotNull(received);
assertThat(received.getHeaders().get(FileHeaders.REMOTE_FILE_INFO), instanceOf(FtpFileInfo.class));
this.adapter.stop();
assertThat(received.getHeaders().get(FileHeaders.REMOTE_FILE_INFO), instanceOf(FtpFileInfo.class));
}
@Test
public void testMaxFetch() {
FtpStreamingMessageSource messageSource = buildsource();
public void testMaxFetch() throws IOException {
FtpStreamingMessageSource messageSource = buildSource();
messageSource.setFilter(new AcceptAllFileListFilter<>());
messageSource.afterPropertiesSet();
Message<InputStream> received = messageSource.receive();
assertNotNull(received);
assertThat(received.getHeaders().get(FileHeaders.REMOTE_FILE), equalTo(" ftpSource1.txt"));
Closeable closeableResource = StaticMessageHeaderAccessor.getCloseableResource(received);
if (closeableResource != null) {
closeableResource.close();
}
}
@Test
public void testMaxFetchNoFilter() {
FtpStreamingMessageSource messageSource = buildsource();
public void testMaxFetchNoFilter() throws IOException {
FtpStreamingMessageSource messageSource = buildSource();
messageSource.setFilter(null);
messageSource.afterPropertiesSet();
Message<InputStream> received = messageSource.receive();
assertNotNull(received);
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(),
Comparator.comparing(FileInfo::getFilename));
messageSource.setRemoteDirectory("ftpSource/");
@@ -169,12 +184,20 @@ public class FtpStreamingMessageSourceTests extends FtpTestSupport {
return pollerMetadata;
}
@Bean
public ConcurrentMap<String, String> metadataMap() {
return new ConcurrentHashMap<>();
}
@Bean
@InboundChannelAdapter(channel = "stream", autoStartup = "false")
public MessageSource<InputStream> ftpMessageSource() {
FtpStreamingMessageSource messageSource = new FtpStreamingMessageSource(template(),
Comparator.comparing(FileInfo::getFilename));
messageSource.setFilter(new AcceptAllFileListFilter<>());
messageSource.setFilter(
new FtpPersistentAcceptOnceFileListFilter(
new SimpleMetadataStore(metadataMap()), "testStreaming"));
messageSource.setRemoteDirectory("ftpSource/");
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.assertThat;
import java.io.IOException;
import java.io.InputStream;
import java.util.Arrays;
import java.util.Comparator;
@@ -59,6 +60,7 @@ import com.jcraft.jsch.ChannelSftp.LsEntry;
/**
* @author Gary Russell
* @author Artem Bilan
*
* @since 4.3
*
*/
@@ -119,7 +121,7 @@ public class SftpStreamingMessageSourceTests extends SftpTestSupport {
}
@Test
public void testMaxFetch() {
public void testMaxFetch() throws IOException {
SftpStreamingMessageSource messageSource = buildSource();
messageSource.setFilter(new AcceptAllFileListFilter<>());
messageSource.afterPropertiesSet();
@@ -127,10 +129,12 @@ public class SftpStreamingMessageSourceTests extends SftpTestSupport {
assertNotNull(received);
assertThat(received.getHeaders().get(FileHeaders.REMOTE_FILE),
anyOf(equalTo(" sftpSource1.txt"), equalTo("sftpSource2.txt")));
received.getPayload().close();
}
@Test
public void testMaxFetchNoFilter() {
public void testMaxFetchNoFilter() throws IOException {
SftpStreamingMessageSource messageSource = buildSource();
messageSource.setFilter(null);
messageSource.afterPropertiesSet();
@@ -138,10 +142,12 @@ public class SftpStreamingMessageSourceTests extends SftpTestSupport {
assertNotNull(received);
assertThat(received.getHeaders().get(FileHeaders.REMOTE_FILE),
anyOf(equalTo(" sftpSource1.txt"), equalTo("sftpSource2.txt")));
received.getPayload().close();
}
@Test
public void testMaxFetchLambdaFilter() {
public void testMaxFetchLambdaFilter() throws IOException {
SftpStreamingMessageSource messageSource = buildSource();
messageSource.setFilter(f -> Arrays.asList(f));
messageSource.afterPropertiesSet();
@@ -149,6 +155,8 @@ public class SftpStreamingMessageSourceTests extends SftpTestSupport {
assertNotNull(received);
assertThat(received.getHeaders().get(FileHeaders.REMOTE_FILE),
anyOf(equalTo(" sftpSource1.txt"), equalTo("sftpSource2.txt")));
received.getPayload().close();
}
private SftpStreamingMessageSource buildSource() {