diff --git a/clearing-parent/dbf-importer/pom.xml b/clearing-parent/dbf-importer/pom.xml index b6275f413..ad2e3f7bc 100644 --- a/clearing-parent/dbf-importer/pom.xml +++ b/clearing-parent/dbf-importer/pom.xml @@ -87,11 +87,11 @@ - - - - - + + org.springframework.boot + spring-boot-starter-test + test + org.junit.jupiter junit-jupiter diff --git a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/config/SFTPConfig.java b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/config/SFTPConfig.java index 274f69025..593108f98 100644 --- a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/config/SFTPConfig.java +++ b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/config/SFTPConfig.java @@ -4,29 +4,27 @@ import com.jcraft.jsch.ChannelSftp; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; -import org.springframework.integration.annotation.InboundChannelAdapter; +import org.springframework.integration.annotation.Gateway; +import org.springframework.integration.annotation.MessagingGateway; import org.springframework.integration.annotation.ServiceActivator; -import org.springframework.integration.core.MessageSource; -import org.springframework.integration.file.filters.AcceptAllFileListFilter; +import org.springframework.integration.channel.DirectChannel; +import org.springframework.integration.file.remote.gateway.AbstractRemoteFileOutboundGateway; import org.springframework.integration.file.remote.session.CachingSessionFactory; import org.springframework.integration.file.remote.session.SessionFactory; -import org.springframework.integration.scheduling.PollerMetadata; -import org.springframework.integration.sftp.inbound.SftpInboundFileSynchronizer; -import org.springframework.integration.sftp.inbound.SftpInboundFileSynchronizingMessageSource; +import org.springframework.integration.sftp.gateway.SftpOutboundGateway; import org.springframework.integration.sftp.session.DefaultSftpSessionFactory; +import org.springframework.messaging.MessageChannel; import org.springframework.messaging.MessageHandler; -import org.springframework.scheduling.support.PeriodicTrigger; import ru.spcex.clearing.dbf.importer.config.settings.ImportDBFServiceSettings; import java.io.File; -import java.util.Comparator; -import java.util.concurrent.TimeUnit; +import java.util.List; + +import static org.springframework.integration.file.remote.gateway.AbstractRemoteFileOutboundGateway.Command.MGET; @Configuration -@ConditionalOnProperty(value="import-dbf-service.store.sftp-in.sftp-src-dir") public class SFTPConfig { private final Logger log = LoggerFactory.getLogger(getClass()); private final ImportDBFServiceSettings settings; @@ -47,40 +45,26 @@ public class SFTPConfig { return new CachingSessionFactory<>(factory); } - @Bean(name = PollerMetadata.DEFAULT_POLLER) - public PollerMetadata defaultPoller() { - PollerMetadata pollerMetadata = new PollerMetadata(); - pollerMetadata.setTrigger(new PeriodicTrigger(5, TimeUnit.SECONDS)); - return pollerMetadata; + @MessagingGateway + public interface DbfGateway { + @Gateway(requestChannel = "listSftpChannel") + List listFiles(String dir); } @Bean - SftpInboundFileSynchronizer sftpInboundFileSynchronizer() { - SftpInboundFileSynchronizer fileSync = new SftpInboundFileSynchronizer(sftpSessionFactory()); - fileSync.setDeleteRemoteFiles(true); - fileSync.setTemporaryFileSuffix(".tmp"); - fileSync.setRemoteDirectory(settings.getStore().getSftpIn().getSftpSrcDir()); - fileSync.setFilter(new AcceptAllFileListFilter<>()); - fileSync.setPreserveTimestamp(true); - fileSync.setComparator(Comparator.comparingInt(f -> f.getAttrs().getMTime())); - return fileSync; + public MessageChannel listSftpChannel(SessionFactory sessionFactory, ImportDBFServiceSettings settings) { + DirectChannel dc = new DirectChannel(); + dc.subscribe(handlerList(sessionFactory, settings)); + return dc; } @Bean - @InboundChannelAdapter("sftpChannel") - public MessageSource sftpMessageSource() { - SftpInboundFileSynchronizingMessageSource source = new SftpInboundFileSynchronizingMessageSource(sftpInboundFileSynchronizer()); - source.setLocalDirectory(new File(settings.getStore().getSrcDir())); - source.setAutoCreateLocalDirectory(true); - return source; - } - - @Bean - @ServiceActivator(inputChannel="sftpChannel") - MessageHandler messageHandler() { - return arg0 -> { - File f = (File) arg0.getPayload(); - log.info("copy file '{}' from sftp event", f.getName()); - }; + @ServiceActivator(inputChannel = "listSftpChannel") + public MessageHandler handlerList(SessionFactory sessionFactory, ImportDBFServiceSettings settings) { + SftpOutboundGateway sftpOutboundGateway = new SftpOutboundGateway(sessionFactory, MGET.getCommand(), null); + sftpOutboundGateway.setLocalDirectory(new File(settings.getStore().getSrcDir())); + sftpOutboundGateway.setAutoCreateLocalDirectory(true); + sftpOutboundGateway.setOption(AbstractRemoteFileOutboundGateway.Option.DELETE); + return sftpOutboundGateway; } } diff --git a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/DBFImporterService.java b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/DBFImporterService.java index 111a222aa..40a243e9b 100644 --- a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/DBFImporterService.java +++ b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/DBFImporterService.java @@ -48,6 +48,7 @@ public class DBFImporterService { } else { log.info("adding import task {}", specificTable); } + fileChecker.checkAndLoadSFTP(); executorService.execute(() -> { Map> newFiles = null; try { diff --git a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/FileChecker.java b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/FileChecker.java index b9e3d74e0..4fadc4fb0 100644 --- a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/FileChecker.java +++ b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/FileChecker.java @@ -1,24 +1,34 @@ package ru.spcex.clearing.dbf.importer.services; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import org.springframework.stereotype.Service; +import ru.spcex.clearing.dbf.importer.config.SFTPConfig; import ru.spcex.clearing.dbf.importer.config.settings.ImportDBFServiceSettings; import ru.spcex.clearing.dbf.importer.logic.data.enums.ETable; +import ru.spcex.platform.utils.log.ExceptionUtils; import java.io.File; +import java.io.IOException; +import java.nio.file.Files; +import java.nio.file.Path; +import java.nio.file.Paths; import java.util.*; +import java.util.stream.Collectors; +import java.util.stream.Stream; @Service("fileChecker") public class FileChecker { private final ImportDBFServiceSettings settings; + private final Logger log = LoggerFactory.getLogger(this.getClass()); + private final SFTPConfig.DbfGateway gateway; - public FileChecker(ImportDBFServiceSettings settings) { + public FileChecker(SFTPConfig.DbfGateway gateway, + ImportDBFServiceSettings settings) { + this.gateway = gateway; this.settings = settings; } - public Map> checkNewFiles() { - return checkNewFiles(null); - } - public Map> checkNewFiles(ETable specificTable) { Map> newFiles = new EnumMap<>(ETable.class); @@ -39,19 +49,30 @@ public class FileChecker { } private List lsDBF(String dbfDirPath) { - File dbfDir = new File(dbfDirPath); - File[] dbfFiles = dbfDir.listFiles((dir, name) -> { - int formatPosition = name.lastIndexOf("."); - if (formatPosition == -1 || formatPosition == name.length() - 1) return false; - return "dbf".equalsIgnoreCase(name.substring(formatPosition + 1)); - }); - List resultFiles = new LinkedList<>(); - if (dbfFiles != null && dbfFiles.length >= 1) { - Arrays.sort(dbfFiles, Comparator.comparingLong(File::lastModified)); - resultFiles.addAll(Arrays.asList(dbfFiles)); + try (Stream stream = Files.walk(Paths.get(dbfDirPath))) { + resultFiles = stream + .filter(file -> !Files.isDirectory(file)) + .map(Path::toFile) + .filter(file -> file.getName().toLowerCase(Locale.ROOT).endsWith(".dbf")) + .sorted(Comparator.comparingLong(File::lastModified)) + .collect(Collectors.toList()); + } catch (IOException e) { + log.error("Error reading 'service.store.src-dir' : {}", ExceptionUtils.getStackTrace(e)); } return resultFiles; } + + public void checkAndLoadSFTP() { + if (settings.getStore().getSftpIn().getSftpSrcPayValDir() == null || settings.getStore().getSftpIn().getSftpSrcPayValDir().size() == 0) { + log.error("No SFTP scanning directories, need will be adding setting like 'import-dbf-service.store.sftp-in.sftp-src-pay-val-dir.VAL=/VAL' and restart app"); + return; + } + List files = new ArrayList<>(); + for (String path : settings.getStore().getSftpIn().getSftpSrcPayValDir().values()) { + files.addAll(gateway.listFiles(path)); + } + if (files.size() > 0) log.info("Loaded from SFTP files count={}", files.size()); + } } diff --git a/clearing-parent/dbf-importer/src/main/resources/application.properties b/clearing-parent/dbf-importer/src/main/resources/application.properties index d9da26709..7957a0c9e 100644 --- a/clearing-parent/dbf-importer/src/main/resources/application.properties +++ b/clearing-parent/dbf-importer/src/main/resources/application.properties @@ -14,6 +14,8 @@ import-dbf-service.common.insert-batch-size=100 import-dbf-service.common.threads-count=10 #sftpSrcDir +import-dbf-service.store.sftp-in.sftp-src-pay-val-dir.rub=clearing_importer_sftp/rub/ +import-dbf-service.store.sftp-in.sftp-src-pay-val-dir.eur=clearing_importer_sftp/eur/ import-dbf-service.store.sftp-in.sftp-src-dir=clearing_importer_sftp import-dbf-service.store.sftp-in.user=user import-dbf-service.store.sftp-in.password=********* diff --git a/clearing-parent/dbf-importer/src/test/java/ru/spcex/clearing/dbf/importer/logic/data/tables/FileWorkTest.java b/clearing-parent/dbf-importer/src/test/java/ru/spcex/clearing/dbf/importer/logic/data/tables/FileWorkTest.java new file mode 100644 index 000000000..31972e5a0 --- /dev/null +++ b/clearing-parent/dbf-importer/src/test/java/ru/spcex/clearing/dbf/importer/logic/data/tables/FileWorkTest.java @@ -0,0 +1,46 @@ +package ru.spcex.clearing.dbf.importer.logic.data.tables; + +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.Disabled; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.context.properties.EnableConfigurationProperties; +import org.springframework.test.context.ContextConfiguration; +import org.springframework.test.context.junit.jupiter.SpringExtension; +import ru.spcex.clearing.dbf.importer.config.settings.ImportDBFServiceSettings; +import ru.spcex.clearing.dbf.importer.logic.data.enums.ETable; +import ru.spcex.clearing.dbf.importer.logic.stages.ChangeDirOfFileStage; +import ru.spcex.clearing.dbf.importer.services.FileChecker; + +import java.io.File; +import java.nio.file.Path; +import java.nio.file.Paths; +import java.util.List; +import java.util.Map; + +@ContextConfiguration(classes = {FileChecker.class, ChangeDirOfFileStage.class}) +@EnableConfigurationProperties(value = ImportDBFServiceSettings.class) +@ExtendWith(SpringExtension.class) +class FileWorkTest { + + @Autowired + FileChecker fileChecker; + + @Autowired + ChangeDirOfFileStage changeDirOfFileStage; + + @BeforeAll + static void setProperty() { + Path path = Paths.get("src", "main", "resources"); + String currentPath = path.toAbsolutePath().toString(); + System.setProperty("spring.config.location", currentPath); + } + + @Disabled + @Test + void checkNewFiles() { + Map> newFiles = fileChecker.checkNewFiles(null); + newFiles.size(); + } +} \ No newline at end of file diff --git a/platform-parent/platform-utils/src/main/java/ru/spcex/platform/utils/config/SftpInboundFolderSetting.java b/platform-parent/platform-utils/src/main/java/ru/spcex/platform/utils/config/SftpInboundFolderSetting.java index 706cf023f..675acec12 100644 --- a/platform-parent/platform-utils/src/main/java/ru/spcex/platform/utils/config/SftpInboundFolderSetting.java +++ b/platform-parent/platform-utils/src/main/java/ru/spcex/platform/utils/config/SftpInboundFolderSetting.java @@ -1,11 +1,14 @@ package ru.spcex.platform.utils.config; +import java.util.Map; + public class SftpInboundFolderSetting { private String sftpSrcDir; private String user; private String password; private String serverIp; private int serverPort; + private Map sftpSrcPayValDir; public String getSftpSrcDir() { return sftpSrcDir; @@ -46,4 +49,12 @@ public class SftpInboundFolderSetting { public void setServerPort(int serverPort) { this.serverPort = serverPort; } + + public Map getSftpSrcPayValDir() { + return sftpSrcPayValDir; + } + + public void setSftpSrcPayValDir(Map sftpSrcPayValDir) { + this.sftpSrcPayValDir = sftpSrcPayValDir; + } }