GH-8691: Add (S)FTP, SMB aged file filters

Fixes https://github.com/spring-projects/spring-integration/issues/8691

* Remove setAge with TimeUnit
Turned out we don't need anymore since we're using Duration for age in FtpLastModifiedFileListFilter and SftpLastModifiedFileListFilter.
* Add changes to the docs
* Introduce AbstractLastModifiedFileListFilter
* Some code readability improvements
* Make language in the docs more official
This commit is contained in:
Adama Sorho
2023-09-02 00:56:31 -04:00
committed by Artem Bilan
parent f3d0441a38
commit 73ed3eeebd
14 changed files with 540 additions and 90 deletions

View File

@@ -930,6 +930,8 @@ project('spring-integration-sftp') {
testImplementation project(':spring-integration-event')
testImplementation project(':spring-integration-file').sourceSets.test.output
testRuntimeOnly 'net.i2p.crypto:eddsa:0.3.0'
}
}

View File

@@ -0,0 +1,128 @@
/*
* Copyright 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.
* You may obtain a copy of the License at
*
* https://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.integration.file.filters;
import java.time.Duration;
import java.time.Instant;
import java.util.ArrayList;
import java.util.List;
import java.util.function.Consumer;
import org.springframework.lang.Nullable;
/**
* The {@link FileListFilter} implementation to filter those files which
* lastModified is less than the {@link #age} in comparison
* with the {@link Instant#now()}.
* When {@link #discardCallback} is provided, it called for all the rejected files.
*
* @param <F> the file
*
* @author Adama Sorho
* @author Artem Bilan
*
* @since 6.2
*/
public abstract class AbstractLastModifiedFileListFilter<F> implements DiscardAwareFileListFilter<F> {
protected static final long ONE_SECOND = 1000;
private static final long DEFAULT_AGE = 60;
private Duration age = Duration.ofSeconds(DEFAULT_AGE);
@Nullable
private Consumer<F> discardCallback;
public AbstractLastModifiedFileListFilter() {
}
public AbstractLastModifiedFileListFilter(Duration age) {
this.age = age;
}
/**
* Set the age that files have to be before being passed by this filter.
* If lastModified plus {@link #age} is before the {@link Instant#now()}, the file
* is filtered.
* Defaults to 60 seconds.
* @param age the Duration.
*/
public void setAge(Duration age) {
this.age = age;
}
/**
* Set the age that files have to be before being passed by this filter.
* If lastModified plus {@link #age} is before the {@link Instant#now()}, the file
* is filtered.
* Defaults to 60 seconds.
* @param age the age in seconds.
*/
public void setAge(long age) {
setAge(Duration.ofSeconds(age));
}
@Override
public void addDiscardCallback(@Nullable Consumer<F> discardCallback) {
this.discardCallback = discardCallback;
}
@Override
public List<F> filterFiles(F[] files) {
List<F> list = new ArrayList<>();
Instant now = Instant.now();
for (F file: files) {
if (fileIsAged(file, now)) {
list.add(file);
}
else if (this.discardCallback != null) {
this.discardCallback.accept(file);
}
}
return list;
}
@Override
public boolean accept(F file) {
if (fileIsAged(file, Instant.now())) {
return true;
}
else if (this.discardCallback != null) {
this.discardCallback.accept(file);
}
return false;
}
private boolean fileIsAged(F file, Instant now) {
return getLastModified(file).plus(this.age).isBefore(now);
}
@Override
public boolean supportsSingleFileFiltering() {
return true;
}
protected Duration getAgeDuration() {
return this.age;
}
protected abstract Instant getLastModified(F remoteFile);
}

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2015-2022 the original author or authors.
* Copyright 2015-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.
@@ -18,51 +18,41 @@ package org.springframework.integration.file.filters;
import java.io.File;
import java.time.Duration;
import java.util.ArrayList;
import java.util.List;
import java.time.Instant;
import java.util.concurrent.TimeUnit;
import java.util.function.Consumer;
import org.springframework.lang.Nullable;
/**
* The {@link FileListFilter} implementation to filter those files which
* {@link File#lastModified()} is less than the {@link #age} in comparison
* {@link File#lastModified()} is less than the age in comparison
* with the current time.
* <p>
* The resolution is done in seconds.
* <p>
* When {@link #discardCallback} is provided, it called for all the
* When discardCallback {@link #addDiscardCallback(Consumer)} is provided, it called for all the
* rejected files.
*
* @author Gary Russell
* @author Artem Bilan
* @author Adama Sorho
*
* @since 4.2
*
*/
public class LastModifiedFileListFilter implements DiscardAwareFileListFilter<File> {
private static final long ONE_SECOND = 1000;
private static final long DEFAULT_AGE = 60;
private volatile long age = DEFAULT_AGE;
@Nullable
private Consumer<File> discardCallback;
public class LastModifiedFileListFilter extends AbstractLastModifiedFileListFilter<File> {
public LastModifiedFileListFilter() {
super();
}
/**
* Construct a {@link LastModifiedFileListFilter} instance with provided {@link #age}.
* Construct a {@link LastModifiedFileListFilter} instance with provided age.
* Defaults to 60 seconds.
* @param age the age in seconds.
* @since 5.0
*/
public LastModifiedFileListFilter(long age) {
this.age = age;
super(Duration.ofSeconds(age));
}
/**
@@ -72,76 +62,25 @@ public class LastModifiedFileListFilter implements DiscardAwareFileListFilter<Fi
* Defaults to 60 seconds.
* @param age the age
* @param unit the timeUnit.
* @deprecated since 6.2 in favor of {@link #setAge(Duration)}
*/
@Deprecated(since = "6.2", forRemoval = true)
public void setAge(long age, TimeUnit unit) {
this.age = unit.toSeconds(age);
setAge(unit.toSeconds(age));
}
/**
* Set the age that files have to be before being passed by this filter.
* If {@link File#lastModified()} plus age is greater than the current time, the file
* is filtered. The resolution is seconds.
* Defaults to 60 seconds.
* @param age the age
* @since 5.1.3
* @return the age in seconds.
* @deprecated since 6.2 in favor of {@link #getAgeDuration()}
*/
public void setAge(Duration age) {
setAge(age.getSeconds());
}
/**
* Set the age that files have to be before being passed by this filter.
* If {@link File#lastModified()} plus age is greater than the current time, the file
* is filtered. The resolution is seconds.
* Defaults to 60 seconds.
* @param age the age
*/
public void setAge(long age) {
setAge(age, TimeUnit.SECONDS);
}
@Deprecated(since = "6.2", forRemoval = true)
public long getAge() {
return this.age;
return getAgeDuration().getSeconds();
}
@Override
public void addDiscardCallback(@Nullable Consumer<File> discardCallbackToSet) {
this.discardCallback = discardCallbackToSet;
}
@Override
public List<File> filterFiles(File[] files) {
List<File> list = new ArrayList<>();
long now = System.currentTimeMillis() / ONE_SECOND;
for (File file : files) {
if (fileIsAged(file, now)) {
list.add(file);
}
else if (this.discardCallback != null) {
this.discardCallback.accept(file);
}
}
return list;
}
@Override
public boolean accept(File file) {
if (fileIsAged(file, System.currentTimeMillis() / ONE_SECOND)) {
return true;
}
else if (this.discardCallback != null) {
this.discardCallback.accept(file);
}
return false;
}
private boolean fileIsAged(File file, long now) {
return file.lastModified() / ONE_SECOND + this.age <= now;
}
@Override
public boolean supportsSingleFileFiltering() {
return true;
protected Instant getLastModified(File file) {
return Instant.ofEpochSecond(file.lastModified() / ONE_SECOND);
}
}

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2015-2022 the original author or authors.
* Copyright 2015-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.
@@ -18,37 +18,38 @@ package org.springframework.integration.file.filters;
import java.io.File;
import java.io.FileOutputStream;
import java.util.concurrent.TimeUnit;
import java.time.Duration;
import java.time.Instant;
import org.junit.Rule;
import org.junit.Test;
import org.junit.rules.TemporaryFolder;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import static org.assertj.core.api.Assertions.assertThat;
/**
* @author Gary Russell
* @author Artem Bilan
* @author Adama Sorho
*
* @since 4.2
*
*/
public class LastModifiedFileListFilterTests {
@Rule
public TemporaryFolder folder = new TemporaryFolder();
@TempDir
public File folder;
@Test
public void testAge() throws Exception {
LastModifiedFileListFilter filter = new LastModifiedFileListFilter();
filter.setAge(60, TimeUnit.SECONDS);
File foo = this.folder.newFile();
filter.setAge(60);
File foo = new File(folder, "test.tmp");
FileOutputStream fileOutputStream = new FileOutputStream(foo);
fileOutputStream.write("x".getBytes());
fileOutputStream.close();
assertThat(filter.filterFiles(new File[] {foo})).hasSize(0);
assertThat(filter.accept(foo)).isFalse();
// Make a file as of yesterday's
foo.setLastModified(System.currentTimeMillis() - 1000 * 60 * 60 * 24);
foo.setLastModified(Instant.now().minus(Duration.ofDays(1)).toEpochMilli());
assertThat(filter.filterFiles(new File[] {foo})).hasSize(1);
assertThat(filter.accept(foo)).isTrue();
}

View File

@@ -0,0 +1,66 @@
/*
* Copyright 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.
* You may obtain a copy of the License at
*
* https://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.integration.ftp.filters;
import java.time.Duration;
import java.time.Instant;
import java.util.function.Consumer;
import org.apache.commons.net.ftp.FTPFile;
import org.springframework.integration.file.filters.AbstractLastModifiedFileListFilter;
/**
* The {@link AbstractLastModifiedFileListFilter} implementation to filter those files which
* {@link FTPFile#getTimestampInstant()} is less than the age in comparison with the {@link Instant#now()}.
* When discardCallback {@link #addDiscardCallback(Consumer)} is provided, it called for all the rejected files.
*
* @author Adama Sorho
* @author Artem Bilan
*
* @since 6.2
*/
public class FtpLastModifiedFileListFilter extends AbstractLastModifiedFileListFilter<FTPFile> {
public FtpLastModifiedFileListFilter() {
super();
}
/**
* Construct a {@link FtpLastModifiedFileListFilter} instance with provided age.
* Defaults to 60 seconds.
* @param age the age in seconds.
*/
public FtpLastModifiedFileListFilter(long age) {
this(Duration.ofSeconds(age));
}
/**
* Construct a {@link FtpLastModifiedFileListFilter} instance with provided age.
* Defaults to 60 seconds.
* @param age the Duration
*/
public FtpLastModifiedFileListFilter(Duration age) {
super(age);
}
@Override
protected Instant getLastModified(FTPFile remoteFile) {
return remoteFile.getTimestampInstant();
}
}

View File

@@ -0,0 +1,55 @@
/*
* Copyright 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.
* You may obtain a copy of the License at
*
* https://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.integration.ftp.filters;
import java.util.Calendar;
import org.apache.commons.net.ftp.FTPFile;
import org.junit.jupiter.api.Test;
import static org.assertj.core.api.Assertions.assertThat;
/**
* @author Adama Sorho
*
* @since 6.2
*/
public class FtpLastModifiedFileListFilterTests {
@Test
public void testAge() {
FtpLastModifiedFileListFilter filter = new FtpLastModifiedFileListFilter();
FTPFile ftpFile1 = new FTPFile();
ftpFile1.setName("foo");
ftpFile1.setTimestamp(Calendar.getInstance());
FTPFile ftpFile2 = new FTPFile();
ftpFile2.setName("bar");
ftpFile2.setTimestamp(Calendar.getInstance());
FTPFile[] files = new FTPFile[] {ftpFile1, ftpFile2};
assertThat(filter.filterFiles(files)).hasSize(0);
assertThat(filter.accept(ftpFile2)).isFalse();
// Make a file as of yesterday's
final Calendar calendar = Calendar.getInstance();
calendar.add(Calendar.DATE, -1);
ftpFile2.setTimestamp(calendar);
assertThat(filter.filterFiles(files)).hasSize(1);
assertThat(filter.accept(ftpFile2)).isTrue();
}
}

View File

@@ -0,0 +1,67 @@
/*
* Copyright 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.
* You may obtain a copy of the License at
*
* https://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.integration.sftp.filters;
import java.nio.file.attribute.FileTime;
import java.time.Duration;
import java.time.Instant;
import java.util.function.Consumer;
import org.apache.sshd.sftp.client.SftpClient;
import org.springframework.integration.file.filters.AbstractLastModifiedFileListFilter;
/**
* The {@link AbstractLastModifiedFileListFilter} implementation to filter those files which
* {@link FileTime#toInstant()} is less than the age in comparison with the {@link Instant#now()}.
* When discardCallback {@link #addDiscardCallback(Consumer)} is provided, it called for all the rejected files.
*
* @author Adama Sorho
* @author Artem Bilan
*
* @since 6.2
*/
public class SftpLastModifiedFileListFilter extends AbstractLastModifiedFileListFilter<SftpClient.DirEntry> {
public SftpLastModifiedFileListFilter() {
super();
}
/**
* Construct a {@link SftpLastModifiedFileListFilter} instance with provided age.
* Defaults to 60 seconds.
* @param age the age in seconds.
*/
public SftpLastModifiedFileListFilter(long age) {
this(Duration.ofSeconds(age));
}
/**
* Construct a {@link SftpLastModifiedFileListFilter} instance with provided age.
* Defaults to 60 seconds.
* @param age the Duration
*/
public SftpLastModifiedFileListFilter(Duration age) {
super(age);
}
@Override
protected Instant getLastModified(SftpClient.DirEntry remoteFile) {
return remoteFile.getAttributes().getModifyTime().toInstant();
}
}

View File

@@ -0,0 +1,57 @@
/*
* Copyright 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.
* You may obtain a copy of the License at
*
* https://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.integration.sftp.filters;
import java.nio.file.attribute.FileTime;
import java.time.Duration;
import java.time.Instant;
import org.apache.sshd.sftp.client.SftpClient;
import org.junit.jupiter.api.Test;
import static org.assertj.core.api.Assertions.assertThat;
/**
* @author Adama Sorho
* @author Artem Bilan
*
* @since 6.2
*/
public class SftpLastModifiedFileListFilterTests {
@Test
public void testAge() {
SftpLastModifiedFileListFilter filter = new SftpLastModifiedFileListFilter();
SftpClient.Attributes attributes1 = new SftpClient.Attributes();
attributes1.setModifyTime(FileTime.from(Instant.now()));
SftpClient.Attributes attributes2 = new SftpClient.Attributes();
attributes2.setModifyTime(FileTime.from(Instant.now()));
SftpClient.DirEntry sftpFile1 = new SftpClient.DirEntry("foo", "foo", attributes1);
SftpClient.DirEntry sftpFile2 = new SftpClient.DirEntry("bar", "bar", attributes2);
SftpClient.DirEntry[] files = new SftpClient.DirEntry[] {sftpFile1, sftpFile2};
assertThat(filter.filterFiles(files)).hasSize(0);
assertThat(filter.accept(sftpFile2)).isFalse();
FileTime fileTime = FileTime.from(Instant.now().minus(Duration.ofDays(1)));
sftpFile2.getAttributes().setModifyTime(fileTime);
assertThat(filter.filterFiles(files)).hasSize(1);
assertThat(filter.accept(sftpFile2)).isTrue();
}
}

View File

@@ -0,0 +1,59 @@
/*
* Copyright 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.
* You may obtain a copy of the License at
*
* https://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.integration.smb.filters;
import java.time.Duration;
import java.time.Instant;
import java.util.function.Consumer;
import jcifs.smb.SmbFile;
import org.springframework.integration.file.filters.AbstractLastModifiedFileListFilter;
/**
* The {@link AbstractLastModifiedFileListFilter} implementation to filter those files which
* {@link SmbFile#getLastModified()} is less than the age in comparison with the current time.
* <p>
* The resolution is done in seconds.
* </p>
* When discardCallback {@link #addDiscardCallback(Consumer)} is provided, it called for all the rejected files.
*
* @author Adama Sorho
*
* @since 6.2
*/
public class SmbLastModifiedFileListFilter extends AbstractLastModifiedFileListFilter<SmbFile> {
public SmbLastModifiedFileListFilter() {
super();
}
/**
* Construct a {@link SmbLastModifiedFileListFilter} instance with provided age.
* Defaults to 60 seconds.
* @param age the age in seconds.
*/
public SmbLastModifiedFileListFilter(long age) {
super(Duration.ofSeconds(age));
}
@Override
protected Instant getLastModified(SmbFile remoteFile) {
return Instant.ofEpochSecond(remoteFile.getLastModified() / ONE_SECOND);
}
}

View File

@@ -0,0 +1,55 @@
/*
* Copyright 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.
* You may obtain a copy of the License at
*
* https://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.integration.smb.filters;
import java.time.Duration;
import java.time.Instant;
import jcifs.smb.SmbFile;
import org.junit.jupiter.api.Test;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.when;
/**
* @author Adama Sorho
* @author Artem Bilan
*
* @since 6.2
*/
public class SmbLastModifiedFileListFilterTests {
@Test
public void testAge() {
SmbLastModifiedFileListFilter filter = new SmbLastModifiedFileListFilter();
filter.setAge(80);
SmbFile smbFile1 = mock(SmbFile.class);
when(smbFile1.getLastModified()).thenReturn(System.currentTimeMillis());
SmbFile smbFile2 = mock(SmbFile.class);
when(smbFile2.getLastModified()).thenReturn(System.currentTimeMillis());
SmbFile smbFile3 = mock(SmbFile.class);
when(smbFile3.getLastModified())
.thenReturn(Instant.now().minus(Duration.ofDays(1)).toEpochMilli());
SmbFile[] files = new SmbFile[] {smbFile1, smbFile2, smbFile3};
assertThat(filter.filterFiles(files)).hasSize(1);
assertThat(filter.accept(smbFile1)).isFalse();
assertThat(filter.accept(smbFile3)).isTrue();
}
}

View File

@@ -98,6 +98,11 @@ You should also understand that the FTP inbound channel adapter is a polling con
Therefore, you must configure a poller (by using either a global default or a local sub-element).
Once a file has been transferred, a message with a `java.io.File` as its payload is generated and sent to the channel identified by the `channel` attribute.
Starting with version 6.2, you can filter FTP files based on last-modified strategy using `FtpLastModifiedFileListFilter`.
This filter can be configured with an `age` property so that only files older than this value are passed by the filter.
The age defaults to 60 seconds, but you should choose an age that is large enough to avoid picking up a file early (due to, say, network glitches).
Look into its Javadoc for more information.
[[more-on-file-filtering-and-incomplete-files]]
== More on File Filtering and Incomplete Files

View File

@@ -95,6 +95,11 @@ SFTP inbound channel adapter is a polling consumer.
Therefore, you must configure a poller (either a global default or a local element).
Once the file has been transferred to a local directory, a message with `java.io.File` as its payload type is generated and sent to the channel identified by the `channel` attribute.
Starting with version 6.2, you can filter SFTP files based on last-modified strategy using `SftpLastModifiedFileListFilter`.
This filter can be configured with an `age` property so that only files older than this value are passed by the filter.
The age defaults to 60 seconds, but you should choose an age that is large enough to avoid picking up a file early (due to, say, network glitches).
Look into its Javadoc for more information.
[[more-on-file-filtering-and-large-files]]
== More on File Filtering and Large Files

View File

@@ -140,6 +140,11 @@ public MessageSource<File> smbMessageSource() {
For XML configuration the `<int-smb:inbound-channel-adapter>` component is provided.
Starting with version 6.2, you can filter SMB files based on last-modified strategy using `SmbLastModifiedFileListFilter`.
This filter can be configured with an `age` property so that only files older than this value are passed by the filter.
The age defaults to 60 seconds, but you should choose an age that is large enough to avoid picking up a file early (due to, say, network glitches.
Look into its Javadoc for more information.
[[configuring-with-the-java-dsl]]
=== Configuring with the Java DSL

View File

@@ -66,3 +66,9 @@ See xref:jdbc/message-store.adoc#jdbc-db-init[Initializing the Database] for mor
A new option `setCreateIndexes(boolean)` has been introduced in `AbstractConfigurableMongoDbMessageStore` to disable the auto indexes creation.
See xref:mongodb.adoc#mongodb-message-store[MongoDB Message Store] for an example.
[[x6.2-remote-files]]
=== Remote Files Support Changes
`FtpLastModifiedFileListFilter`, `SftpLastModifiedFileListFilter` and `SmbLastModifiedFileListFilter` have been introduced to allow files filtering based on a last-modified strategy respectively for `FTP`, `SFTP` and `SMB`.
See xref:ftp/inbound.adoc#ftp-inbound[FTP Inbound Channel Adapter], xref:sftp/inbound.adoc#sftp-inbound[SFTP Inbound Channel Adapter], and xref:smb.adoc#smb-inbound[SMB Inbound Channel Adapter] for more information.