diff --git a/spring-integration-file/src/main/java/org/springframework/integration/file/remote/AbstractRemoteFileStreamingMessageSource.java b/spring-integration-file/src/main/java/org/springframework/integration/file/remote/AbstractRemoteFileStreamingMessageSource.java index d267a2352d..a3793f9e2d 100644 --- a/spring-integration-file/src/main/java/org/springframework/integration/file/remote/AbstractRemoteFileStreamingMessageSource.java +++ b/spring-integration-file/src/main/java/org/springframework/integration/file/remote/AbstractRemoteFileStreamingMessageSource.java @@ -215,7 +215,7 @@ public abstract class AbstractRemoteFileStreamingMessageSource if (this.filter != null && this.filter.supportsSingleFileFiltering() && !this.filter.accept(file.getFileInfo())) { - if (this.toBeReceived.size() > 0) { // don't re-fetch already filtered files + if (!this.toBeReceived.isEmpty()) { // don't re-fetch already filtered files file = poll(); continue; } @@ -267,7 +267,7 @@ public abstract class AbstractRemoteFileStreamingMessageSource } protected AbstractFileInfo poll() { - if (this.toBeReceived.size() == 0) { + if (this.toBeReceived.isEmpty()) { listFiles(); } return this.toBeReceived.poll(); @@ -297,7 +297,7 @@ public abstract class AbstractRemoteFileStreamingMessageSource if (!ObjectUtils.isEmpty(files)) { List> fileInfoList; if (this.filter != null && !this.filter.supportsSingleFileFiltering()) { - int maxFetchSize = getMaxFetchSize(); + int maxFetchSize = getMaxFetchSize() - this.fetched.get(); List filteredFiles = this.filter.filterFiles(files); if (maxFetchSize > 0 && filteredFiles.size() > maxFetchSize) { rollbackFromFileToListEnd(filteredFiles, filteredFiles.get(maxFetchSize)); 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 6f026a0171..0352691779 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 @@ -1,5 +1,5 @@ /* - * Copyright 2016-2022 the original author or authors. + * Copyright 2016-2023 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -16,13 +16,17 @@ package org.springframework.integration.sftp.inbound; +import java.io.File; +import java.io.IOException; import java.io.InputStream; +import java.nio.charset.StandardCharsets; import java.time.Duration; import java.util.Arrays; import java.util.Comparator; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; +import org.apache.commons.io.FileUtils; import org.apache.sshd.sftp.client.SftpClient; import org.junit.jupiter.api.Test; @@ -39,11 +43,14 @@ import org.springframework.integration.core.MessageSource; import org.springframework.integration.endpoint.SourcePollingChannelAdapter; import org.springframework.integration.file.FileHeaders; import org.springframework.integration.file.filters.AcceptAllFileListFilter; +import org.springframework.integration.file.filters.ChainFileListFilter; import org.springframework.integration.file.remote.session.SessionFactory; import org.springframework.integration.metadata.SimpleMetadataStore; import org.springframework.integration.scheduling.PollerMetadata; import org.springframework.integration.sftp.SftpTestSupport; import org.springframework.integration.sftp.filters.SftpPersistentAcceptOnceFileListFilter; +import org.springframework.integration.sftp.filters.SftpSimplePatternFileListFilter; +import org.springframework.integration.sftp.filters.SftpSystemMarkerFilePresentFileListFilter; import org.springframework.integration.sftp.session.SftpFileInfo; import org.springframework.integration.sftp.session.SftpRemoteFileTemplate; import org.springframework.integration.transformer.StreamTransformer; @@ -168,6 +175,69 @@ public class SftpStreamingMessageSourceTests extends SftpTestSupport { StaticMessageHeaderAccessor.getCloseableResource(received).close(); } + + @Test + public void maxFetchIsAdjustedWhenNoSupportsSingleFileFiltering() throws Exception { + SftpStreamingMessageSource messageSource = buildSource(); + ChainFileListFilter chainFileListFilter = new ChainFileListFilter<>(); + SftpSystemMarkerFilePresentFileListFilter sftpSystemMarkerFilePresentFileListFilter = + new SftpSystemMarkerFilePresentFileListFilter( + new SftpSimplePatternFileListFilter("*"), ".trg"); + SftpPersistentAcceptOnceFileListFilter sftpPersistentAcceptOnceFileListFilter = + new SftpPersistentAcceptOnceFileListFilter(new SimpleMetadataStore(), "prefix"); + chainFileListFilter.addFilter(sftpSystemMarkerFilePresentFileListFilter); + chainFileListFilter.addFilter(sftpPersistentAcceptOnceFileListFilter); + messageSource.setFilter(chainFileListFilter); + messageSource.setMaxFetchSize(5); + messageSource.afterPropertiesSet(); + messageSource.start(); + + addFileAndTrigger("file001"); + addFileAndTrigger("file002"); + + Message received = messageSource.receive(); + assertThat(received).isNotNull(); + assertThat(received.getHeaders().get(FileHeaders.REMOTE_FILE)).isEqualTo("file001"); + + received = messageSource.receive(); + assertThat(received).isNotNull(); + assertThat(received.getHeaders().get(FileHeaders.REMOTE_FILE)).isEqualTo("file002"); + + addFileAndTrigger("file003"); + addFileAndTrigger("file004"); + addFileAndTrigger("file005"); + addFileAndTrigger("file006"); + addFileAndTrigger("file007"); + + received = messageSource.receive(); + assertThat(received).isNotNull(); + assertThat(received.getHeaders().get(FileHeaders.REMOTE_FILE)).isEqualTo("file003"); + + received = messageSource.receive(); + assertThat(received).isNotNull(); + assertThat(received.getHeaders().get(FileHeaders.REMOTE_FILE)).isEqualTo("file004"); + + received = messageSource.receive(); + assertThat(received).isNotNull(); + assertThat(received.getHeaders().get(FileHeaders.REMOTE_FILE)).isEqualTo("file005"); + + received = messageSource.receive(); + assertThat(received).isNotNull(); + assertThat(received.getHeaders().get(FileHeaders.REMOTE_FILE)).isEqualTo("file006"); + + received = messageSource.receive(); + assertThat(received).isNotNull(); + assertThat(received.getHeaders().get(FileHeaders.REMOTE_FILE)).isEqualTo("file007"); + } + + private void addFileAndTrigger(String filename) throws IOException { + File file = new File(this.sourceRemoteDirectory, filename); + FileUtils.writeStringToFile(file, "source1", StandardCharsets.UTF_8); + + file = new File(this.sourceRemoteDirectory, filename + ".trg"); + file.createNewFile(); + } + private SftpStreamingMessageSource buildSource() { SftpStreamingMessageSource messageSource = new SftpStreamingMessageSource(this.config.template(),