From 6ecd948a32b19a49939ef4a04fdf9b84ba76eaba Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Fri, 29 Jun 2018 21:28:28 -0400 Subject: [PATCH] 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** --- .../FtpStreamingMessageSourceTests.java | 51 ++++++++++++++----- .../SftpStreamingMessageSourceTests.java | 14 +++-- 2 files changed, 48 insertions(+), 17 deletions(-) diff --git a/spring-integration-ftp/src/test/java/org/springframework/integration/ftp/inbound/FtpStreamingMessageSourceTests.java b/spring-integration-ftp/src/test/java/org/springframework/integration/ftp/inbound/FtpStreamingMessageSourceTests.java index dc416c3323..5b41504519 100644 --- a/spring-integration-ftp/src/test/java/org/springframework/integration/ftp/inbound/FtpStreamingMessageSourceTests.java +++ b/spring-integration-ftp/src/test/java/org/springframework/integration/ftp/inbound/FtpStreamingMessageSourceTests.java @@ -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 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) 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 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 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 metadataMap() { + return new ConcurrentHashMap<>(); + } + @Bean @InboundChannelAdapter(channel = "stream", autoStartup = "false") public MessageSource 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; } diff --git a/spring-integration-sftp/src/test/java/org/springframework/integration/sftp/inbound/SftpStreamingMessageSourceTests.java b/spring-integration-sftp/src/test/java/org/springframework/integration/sftp/inbound/SftpStreamingMessageSourceTests.java index 9854a3aa98..17d55b76eb 100644 --- a/spring-integration-sftp/src/test/java/org/springframework/integration/sftp/inbound/SftpStreamingMessageSourceTests.java +++ b/spring-integration-sftp/src/test/java/org/springframework/integration/sftp/inbound/SftpStreamingMessageSourceTests.java @@ -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() {