From d3a9c462344e48e9fb25ab7337a3e2c955098031 Mon Sep 17 00:00:00 2001 From: AKurakin Date: Wed, 24 May 2023 20:14:00 +0300 Subject: [PATCH] =?UTF-8?q?swt-exporter=20http://jira.mfd.msk:8088/browse/?= =?UTF-8?q?CLS-317=20=D0=BD=D0=B0=D1=87=D0=B0=D0=BB=D0=BE?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- clearing-parent/pom.xml | 1 + clearing-parent/swt-exporter/pom.xml | 112 ++++++++++++++++++ .../swt/exporter/SwtExportApplication.java | 18 +++ .../swt/exporter/config/ImdgConfig.java | 44 +++++++ .../swt/exporter/config/KafkaConfig.java | 56 +++++++++ .../exporter/config/SwtExporterConfig.java | 11 ++ .../swt/exporter/config/settings/Common.java | 32 +++++ .../settings/ExportSwtServiceSettings.java | 59 +++++++++ .../swt/exporter/config/settings/Store.java | 50 ++++++++ .../services/AbstractExporterService.java | 78 ++++++++++++ .../swt/exporter/services/FileStorage.java | 35 ++++++ .../services/LauncherCommandReceiver.java | 38 ++++++ .../exportimpl/MoneyExporterService.java | 103 ++++++++++++++++ .../src/main/resources/application.properties | 25 ++++ .../src/main/resources/logback.xml | 37 ++++++ .../swt/exporter/AbstractServiceTest.java | 53 +++++++++ .../services/AbstractExporterServiceTest.java | 23 ++++ .../services/MoneyExporterServiceTest.java | 102 ++++++++++++++++ 18 files changed, 877 insertions(+) create mode 100644 clearing-parent/swt-exporter/pom.xml create mode 100644 clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/SwtExportApplication.java create mode 100644 clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/config/ImdgConfig.java create mode 100644 clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/config/KafkaConfig.java create mode 100644 clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/config/SwtExporterConfig.java create mode 100644 clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/config/settings/Common.java create mode 100644 clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/config/settings/ExportSwtServiceSettings.java create mode 100644 clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/config/settings/Store.java create mode 100644 clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/services/AbstractExporterService.java create mode 100644 clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/services/FileStorage.java create mode 100644 clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/services/LauncherCommandReceiver.java create mode 100644 clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/services/exportimpl/MoneyExporterService.java create mode 100644 clearing-parent/swt-exporter/src/main/resources/application.properties create mode 100644 clearing-parent/swt-exporter/src/main/resources/logback.xml create mode 100644 clearing-parent/swt-exporter/src/test/java/ru/spcex/clearing/swt/exporter/AbstractServiceTest.java create mode 100644 clearing-parent/swt-exporter/src/test/java/ru/spcex/clearing/swt/exporter/services/AbstractExporterServiceTest.java create mode 100644 clearing-parent/swt-exporter/src/test/java/ru/spcex/clearing/swt/exporter/services/MoneyExporterServiceTest.java diff --git a/clearing-parent/pom.xml b/clearing-parent/pom.xml index c80817211..b00850e36 100644 --- a/clearing-parent/pom.xml +++ b/clearing-parent/pom.xml @@ -38,6 +38,7 @@ cleaning-builders trade-importer lim-exporter + swt-exporter diff --git a/clearing-parent/swt-exporter/pom.xml b/clearing-parent/swt-exporter/pom.xml new file mode 100644 index 000000000..36f959d2c --- /dev/null +++ b/clearing-parent/swt-exporter/pom.xml @@ -0,0 +1,112 @@ + + 4.0.0 + swt-exporter + Swt exporter + SPCEX-1.0.0.0 + jar + + + ru.spcex.clearing + clearing-parent + SPCEX-1.0.0.0 + + + + 17 + 17 + UTF-8 + + + + + org.springframework.boot + spring-boot-starter + + + com.fasterxml.jackson.core + jackson-databind + + + org.springframework.integration + spring-integration-sftp + + + + ru.spcex.clearing + classes + SPCEX-1.0.0.0 + compile + + + ru.spcex.platform + platform-messaging + + + ru.spcex.platform + platform-imdg-api-hazelcast-impl + + + ru.spcex.platform + platform-enum + + + + + org.springframework.boot + spring-boot-starter-test + test + + + ru.spcex.clearing + test-clearing + test + + + + + + + src/main/resources + + application.properties + + false + + + + + org.springframework.boot + spring-boot-maven-plugin + + + + repackage + + + + + ${project.artifactId} + + + + org.apache.maven.plugins + maven-surefire-plugin + 2.21.0 + + + org.junit.platform + junit-platform-surefire-provider + 1.2.0-M1 + + + org.junit.jupiter + junit-jupiter-engine + 5.2.0-M1 + + + + + + + diff --git a/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/SwtExportApplication.java b/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/SwtExportApplication.java new file mode 100644 index 000000000..6b7de7a40 --- /dev/null +++ b/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/SwtExportApplication.java @@ -0,0 +1,18 @@ +package ru.spcex.clearing.swt.exporter; + +import org.slf4j.LoggerFactory; +import org.springframework.boot.autoconfigure.SpringBootApplication; +import org.springframework.boot.builder.SpringApplicationBuilder; + +@SpringBootApplication +public class SwtExportApplication { + public static void main(String[] args) { + try { + SpringApplicationBuilder builder = new SpringApplicationBuilder(SwtExportApplication.class); + builder.run(args); + } catch (Throwable e) { + LoggerFactory.getLogger(SwtExportApplication.class).error("Swt-exporter start failed: {} -> {}", e.getClass().getSimpleName(), e.getMessage()); + System.exit(-1); + } + } +} diff --git a/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/config/ImdgConfig.java b/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/config/ImdgConfig.java new file mode 100644 index 000000000..0799f35e7 --- /dev/null +++ b/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/config/ImdgConfig.java @@ -0,0 +1,44 @@ +package ru.spcex.clearing.swt.exporter.config; + +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; +import ru.spcex.clearing.swt.exporter.config.settings.ExportSwtServiceSettings; +import ru.spcex.platform.imdg.api.ImdgProvider; +import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService; + +@Configuration +public class ImdgConfig { + private static ThreadPoolTaskExecutor createThreadPoolTaskExecutor(int maxPoolSz, boolean waitForCompletion) { + ThreadPoolTaskExecutor pool = new ThreadPoolTaskExecutor(); + if (maxPoolSz > 2) { + pool.setKeepAliveSeconds(60); + pool.setAllowCoreThreadTimeOut(true); + } + pool.setCorePoolSize(maxPoolSz); + pool.setWaitForTasksToCompleteOnShutdown(waitForCompletion); + return pool; + } + + @Bean(name = "taskExecutorHazelcastClientInitializer") + public ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer() { + return createThreadPoolTaskExecutor(1, true); + } + + @Bean(name = "taskExecutorIdGeneratorAwaiter") + public ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter() { + return createThreadPoolTaskExecutor(1, false); + } + + @Bean("imdgProvider") + public ImdgProvider imdgProvider( + @Qualifier("taskExecutorHazelcastClientInitializer") ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer, + @Qualifier("taskExecutorIdGeneratorAwaiter") ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter, + ExportSwtServiceSettings settings) { + return new HazelcastService(taskExecutorHazelcastClientInitializer, + taskExecutorIdGeneratorAwaiter, + settings.getHazelcast()); + } + +} diff --git a/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/config/KafkaConfig.java b/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/config/KafkaConfig.java new file mode 100644 index 000000000..3ad111793 --- /dev/null +++ b/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/config/KafkaConfig.java @@ -0,0 +1,56 @@ +package ru.spcex.clearing.swt.exporter.config; + +import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.producer.Producer; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.config.ConfigurableBeanFactory; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.Scope; +import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.kafka.core.ProducerFactory; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.swt.exporter.config.settings.ExportSwtServiceSettings; +import ru.spcex.clearing.platform.messaging.config.KafkaConsumerFactory; +import ru.spcex.clearing.platform.messaging.config.KafkaProducerFactory; +import ru.spcex.clearing.platform.messaging.service.RequestInfo; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgId; +import ru.spcex.platform.imdg.api.ImdgProvider; + +@Configuration +public class KafkaConfig { + @Autowired + @Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) + @Bean + public Consumer createConsumer(ExportSwtServiceSettings settings) { + return KafkaConsumerFactory.consumer(settings.getKafkaConsumer()); + } + + @Autowired + @Bean + public Producer createProducer(ExportSwtServiceSettings settings) { + return KafkaProducerFactory.producer(settings.getKafkaProducer()); + } + + @Bean + public KafkaTemplate kafkaTemplate(ProducerFactory pf) { + return new KafkaTemplate<>(pf); + } + + @Autowired + @Bean + public KafkaSender kafkaSender(KafkaTemplate kafkaTemplate, ImdgProvider imdgProvider) { + ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator(); + return KafkaSender + .setup() + .setKafkaTemplate(kafkaTemplate) + .idGenerator(imdgIdGenerator::nextId) + .imdgProvider(s -> { + Imdg imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class); + return imdg::insert; + }) + .build(); + } +} diff --git a/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/config/SwtExporterConfig.java b/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/config/SwtExporterConfig.java new file mode 100644 index 000000000..0f3b81fa2 --- /dev/null +++ b/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/config/SwtExporterConfig.java @@ -0,0 +1,11 @@ +package ru.spcex.clearing.swt.exporter.config; + +import org.springframework.boot.context.properties.EnableConfigurationProperties; +import org.springframework.context.annotation.ComponentScan; +import org.springframework.context.annotation.Configuration; + +@Configuration +@EnableConfigurationProperties +@ComponentScan(basePackages = {"ru.spcex.clearing.swt.exporter"}) +public class SwtExporterConfig { +} diff --git a/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/config/settings/Common.java b/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/config/settings/Common.java new file mode 100644 index 000000000..987a47ff6 --- /dev/null +++ b/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/config/settings/Common.java @@ -0,0 +1,32 @@ +package ru.spcex.clearing.swt.exporter.config.settings; + +public class Common { + + private String encoding; + private int insertBatchSize; + private int threadsCount; + + public String getEncoding() { + return encoding; + } + + public void setEncoding(String encoding) { + this.encoding = encoding; + } + + public int getInsertBatchSize() { + return insertBatchSize; + } + + public void setInsertBatchSize(int insertBatchSize) { + this.insertBatchSize = insertBatchSize; + } + + public int getThreadsCount() { + return threadsCount; + } + + public void setThreadsCount(int threadsCount) { + this.threadsCount = threadsCount; + } +} diff --git a/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/config/settings/ExportSwtServiceSettings.java b/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/config/settings/ExportSwtServiceSettings.java new file mode 100644 index 000000000..625935e2e --- /dev/null +++ b/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/config/settings/ExportSwtServiceSettings.java @@ -0,0 +1,59 @@ +package ru.spcex.clearing.swt.exporter.config.settings; + +import org.springframework.boot.context.properties.ConfigurationProperties; +import org.springframework.context.annotation.PropertySource; +import org.springframework.stereotype.Component; +import ru.spcex.clearing.platform.messaging.config.element.KafkaConsumerSettings; +import ru.spcex.clearing.platform.messaging.config.element.KafkaProducerSettings; +import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams; + +@Component +@PropertySource("file:${spring.config.location}/application.properties") +@ConfigurationProperties("export-swt-service") +public class ExportSwtServiceSettings { + private HazelcastClientParams hazelcast; + private KafkaConsumerSettings kafkaConsumer; + private KafkaProducerSettings kafkaProducer; + private Store sFTPStore; + private Common common; + + public HazelcastClientParams getHazelcast() { + return hazelcast; + } + + public void setHazelcast(HazelcastClientParams hazelcast) { + this.hazelcast = hazelcast; + } + + public KafkaConsumerSettings getKafkaConsumer() { + return kafkaConsumer; + } + + public void setKafkaConsumer(KafkaConsumerSettings kafkaConsumer) { + this.kafkaConsumer = kafkaConsumer; + } + + public Store getStore() { + return sFTPStore; + } + + public void setStore(Store Store) { + this.sFTPStore = Store; + } + + public KafkaProducerSettings getKafkaProducer() { + return kafkaProducer; + } + + public void setKafkaProducer(KafkaProducerSettings kafkaProducer) { + this.kafkaProducer = kafkaProducer; + } + + public Common getCommon() { + return common; + } + + public void setCommon(Common common) { + this.common = common; + } +} diff --git a/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/config/settings/Store.java b/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/config/settings/Store.java new file mode 100644 index 000000000..da0a3e491 --- /dev/null +++ b/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/config/settings/Store.java @@ -0,0 +1,50 @@ +package ru.spcex.clearing.swt.exporter.config.settings; + +public class Store { + + private String outDir; + private String user; + private String password; + private String serverIp; + private int serverPort; + + public String getUser() { + return user; + } + + public void setUser(String user) { + this.user = user; + } + + public String getPassword() { + return password; + } + + public void setPassword(String password) { + this.password = password; + } + + public String getServerIp() { + return serverIp; + } + + public void setServerIp(String serverIp) { + this.serverIp = serverIp; + } + + public int getServerPort() { + return serverPort; + } + + public void setServerPort(int serverPort) { + this.serverPort = serverPort; + } + + public String getOutDir() { + return outDir; + } + + public void setOutDir(String outDir) { + this.outDir = outDir; + } +} diff --git a/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/services/AbstractExporterService.java b/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/services/AbstractExporterService.java new file mode 100644 index 000000000..61e0b1526 --- /dev/null +++ b/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/services/AbstractExporterService.java @@ -0,0 +1,78 @@ +package ru.spcex.clearing.swt.exporter.services; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import ru.clearing.classes.statics.data.registry.Registry; +import ru.clearing.classes.statics.data.registry.TradingClearingRegistry; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.platform.messaging.domain.cud.utilities.LimExportedRequest; +import ru.spcex.clearing.platform.messaging.serialization.LogFormatter; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; + +import java.io.IOException; +import java.nio.file.Files; +import java.time.LocalDateTime; +import java.time.format.DateTimeFormatter; +import java.util.Collection; +import java.util.List; + +import static ru.spcex.clearing.platform.messaging.domain.Consts.LIM_EXPORTED; + +public abstract class AbstractExporterService { + private final Logger log = LoggerFactory.getLogger(getClass()); + protected final Imdg registryImdg; + private final List validStatus = List.of("ACTV", "ROPN"); + private final Imdg tradingClearingRegistryImdg; + private final DateTimeFormatter dtFormatter = DateTimeFormatter.ofPattern("yyyyMMddHHmmss"); + private final KafkaSender kafkaSender; + protected final FileStorage fileStorage; + + protected AbstractExporterService(FileStorage fileStorage, + KafkaSender kafkaSender, ImdgProvider imdgProvider) { + this.fileStorage = fileStorage; + this.kafkaSender = kafkaSender; + this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class); + this.tradingClearingRegistryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class); + } + + public abstract Collection getLimFileRows(); + + public abstract String getTargetFileName(); + + public void process() { + String fileName = getTargetFileName(); + log.debug("Start export {} Lim file", fileName); + + try { + // todo возможная оптимизация: посмотреть размеры файлов, возможно обойтись без временного файла + //Files.write(limFilePath, getLimFileRows()); + fileStorage.saveFile(fileName, null); + } catch (IOException e) { + log.error("Failed export {} file", fileName); + throw new RuntimeException(e); + } + log.debug("Successfully exported {} file", fileName); + + sendSwtxportedNotification(fileName); + } + + void sendSwtxportedNotification(String fileName) { + LimExportedRequest limExportedRequest = new LimExportedRequest(); + limExportedRequest.setLimFileName(fileName); + log.debug("Send message to kafka \"{}\": {}", LIM_EXPORTED, LogFormatter.toStringWrapper(limExportedRequest)); + kafkaSender.sendRequestToQueue(LIM_EXPORTED, limExportedRequest);//todo rewrite!!! + } + + protected String prepareFileName(Long counter, String target) { + String dt = dtFormatter.format(LocalDateTime.now()); + String type="09"; + String section="U"; + String counterS = counter==null?"":"_"+counter; + String partyCode=""; + String result="KS_RDC_DF-%s_%s_PRC%s%s%s.swt".formatted(type, section,dt,counterS,partyCode); + return result; + } + +} diff --git a/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/services/FileStorage.java b/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/services/FileStorage.java new file mode 100644 index 000000000..af9fd848e --- /dev/null +++ b/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/services/FileStorage.java @@ -0,0 +1,35 @@ +package ru.spcex.clearing.swt.exporter.services; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.stereotype.Service; +import ru.spcex.clearing.swt.exporter.config.settings.ExportSwtServiceSettings; + +import java.io.File; +import java.io.IOException; + +@Service +public class FileStorage { + protected final Logger log= LoggerFactory.getLogger(getClass()); + protected File outPath; + + @Autowired + public FileStorage(ExportSwtServiceSettings config) { + if (config.getStore().getOutDir()==null || config.getStore().getOutDir().isBlank()) { + throw new IllegalArgumentException("Out directory settings is empty."); + } + this.outPath = new File(config.getStore().getOutDir()); //todo ... + if (!outPath.isDirectory()) { + log.info("Path not exist. mkdir \"{}\"", outPath.getAbsolutePath()); + if (!outPath.mkdir()) { + log.error("Can not make output directory \"{}\"", outPath); + } + } + log.info("Output directory \"{}\"", outPath); + } + + public void saveFile(String fileName, byte[] data) throws IOException { + //todo ... + } +} diff --git a/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/services/LauncherCommandReceiver.java b/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/services/LauncherCommandReceiver.java new file mode 100644 index 000000000..a3ae54ca6 --- /dev/null +++ b/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/services/LauncherCommandReceiver.java @@ -0,0 +1,38 @@ +package ru.spcex.clearing.swt.exporter.services; + +import org.apache.kafka.clients.consumer.Consumer; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.InitializingBean; +import org.springframework.stereotype.Service; +import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest; +import ru.spcex.clearing.platform.messaging.service.QueueConsumer; +import ru.spcex.clearing.swt.exporter.services.exportimpl.MoneyExporterService; +import ru.spcex.platform.enumeration.Task; + +@Service +public class LauncherCommandReceiver extends QueueConsumer implements InitializingBean { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final MoneyExporterService moneyExporterService; +// private final SecurityExporterService securityExporterService; + + public LauncherCommandReceiver(Consumer kafkaQueue, + MoneyExporterService moneyExporterService + // , SecurityExporterService securityExporterService + ) { + super(kafkaQueue); + this.moneyExporterService = moneyExporterService; +// this.securityExporterService = securityExporterService; + } + + @Override + public void afterPropertiesSet() { + callback(LauncherCommandRequest.class) + .setConsumer(action -> moneyExporterService.process()) + .forDestination(Task.unloadingSession_LIMM.topic(), callbacks::put); // LIMM +// callback(LauncherCommandRequest.class) +// .setConsumer(action -> securityExporterService.process()) +// .forDestination(Task.unloadingSession_LIMS.topic(), callbacks::put); // LIMS + init(); + } +} diff --git a/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/services/exportimpl/MoneyExporterService.java b/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/services/exportimpl/MoneyExporterService.java new file mode 100644 index 000000000..7576ccad1 --- /dev/null +++ b/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/services/exportimpl/MoneyExporterService.java @@ -0,0 +1,103 @@ +package ru.spcex.clearing.swt.exporter.services.exportimpl; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.stereotype.Service; +import ru.clearing.classes.statics.data.registry.Registry; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; +import ru.spcex.clearing.swt.exporter.services.AbstractExporterService; +import ru.spcex.clearing.swt.exporter.services.FileStorage; +import ru.spcex.platform.imdg.api.ImdgProvider; + +import java.math.BigDecimal; +import java.time.LocalDate; +import java.util.ArrayList; +import java.util.Collection; +import java.util.List; +import java.util.Map; +import java.util.stream.Collectors; + +@Service +public class MoneyExporterService extends AbstractExporterService { + private final Logger log = LoggerFactory.getLogger(getClass()); + + public MoneyExporterService(FileStorage fileStorage, + KafkaSender kafkaSender, + ImdgProvider imdgProvider) { + super(fileStorage, kafkaSender, imdgProvider); + } + + @Override + public String getTargetFileName() { + return prepareFileName(null, "money"); + } + + @Override + public Collection getLimFileRows() { + log.debug("Started loading and formation of money file lines"); + LocalDate currentDate = LocalDate.now(); + List swtFileRows = new ArrayList<>(); + Collection registriesA = registryImdg.getCollectionObjectsByFieldValues(Map.of( + "registryDesignation", "A", + "registryInstrumentType", "M", + "registryUnit", "F" + )); + Collection registriesD = registryImdg.getCollectionObjectsByFieldValues(Map.of( + "registryDesignation", "D", + "registryInstrumentType", "M", + "registryUnit", "T" + )); + Map> byTcrA = registriesA.stream() + .collect(Collectors.groupingBy(Registry::getTradingClearingRegistry)); + Map> byTcrD = registriesD.stream() + .collect(Collectors.groupingBy(Registry::getTradingClearingRegistry)); + + for (Map.Entry> entryA : byTcrA.entrySet()) { + List registriesListA = entryA.getValue(); + List registriesListB = byTcrD.get(entryA.getKey()); + for (Registry registryA : registriesListA) { +// todo if (checkNotBlocked(registryA)) { +// Registry registryD = findRegistryBySecurityId(registryA.getSecurityId(), registriesListB); +// limFileRows.add(getRow(registryA, registryD)); +// } + } + } + log.debug("Successfully completed the formation of rows: {} for export money", swtFileRows.size()); + return swtFileRows; + } + + public String getRow(Registry registryA, Registry registryD) { + StringBuilder row = new StringBuilder(); + + row.append("MONEY: FIRM_ID = "); + row.append(registryA.getTradingCode()); + + row.append("; TAG = SPVB"); + + row.append("; CURR_CODE = "); + row.append(registryA.getSecuritySymbol()); + + row.append("; CLIENT_CODE = "); + row.append(registryA.getTradingClearingRegistry()); + + row.append("; OPEN_BALANCE = "); + BigDecimal balance = registryA.getBalance() != null ? + registryD != null && registryD.getBalance() != null ? registryA.getBalance().subtract(registryD.getBalance()) : registryA.getBalance() : + BigDecimal.ZERO; + row.append(balance); + + row.append("; OPEN_LIMIT = 0.00"); + + row.append("; LIMIT_KIND = 0;"); + return row.toString(); + } + + private Registry findRegistryBySecurityId(Long securityId, List registries) { + for (Registry registry : registries) { + if (securityId.equals(registry.getSecurityId())) { + return registry; + } + } + return null; + } +} diff --git a/clearing-parent/swt-exporter/src/main/resources/application.properties b/clearing-parent/swt-exporter/src/main/resources/application.properties new file mode 100644 index 000000000..0106cfc7d --- /dev/null +++ b/clearing-parent/swt-exporter/src/main/resources/application.properties @@ -0,0 +1,25 @@ +spring.main.web-application-type=none + +export-swt-service.hazelcast.cluster-members=10.200.200.181:5701 +export-swt-service.hazelcast.login=dev +export-swt-service.hazelcast.password=dev-pass + +export-swt-service.common.encoding=cp866 +export-swt-service.common.threads-count=10 + +export-swt-service.out-dir=DocOut + +export-swt-service.kafka-consumer.bootstrap-servers=localhost:9092 +export-swt-service.kafka-consumer.group-id=dev-group-balance-service +export-swt-service.kafka-consumer.enable-auto-commit=false +export-swt-service.kafka-consumer.session-timeout-ms=30000 +export-swt-service.kafka-consumer.auto-offset-reset=latest +export-swt-service.kafka-consumer.linger-ms=1 +export-swt-service.kafka-consumer.buffer-memory=33554432 + +export-swt-service.kafka-producer.bootstrap-servers=localhost:9092 +export-swt-service.kafka-producer.acks=all +export-swt-service.kafka-producer.retries=0 +export-swt-service.kafka-producer.batch-size=16384 +export-swt-service.kafka-producer.linger-ms=1 +export-swt-service.kafka-producer.buffer-memory=33554432 \ No newline at end of file diff --git a/clearing-parent/swt-exporter/src/main/resources/logback.xml b/clearing-parent/swt-exporter/src/main/resources/logback.xml new file mode 100644 index 000000000..b88eeb279 --- /dev/null +++ b/clearing-parent/swt-exporter/src/main/resources/logback.xml @@ -0,0 +1,37 @@ + + + + + UTF-8 + %date{HH:mm:ss.SSS} [%thread] %-5level %logger{36} - %msg%n + + + + + ./logs/swt-exporter.log + + UTF-8 + %d{HH:mm:ss.SSS} [%thread] %-5level %logger{36} - %msg%n + + + + ../logs/swt-exporter.%i.log + + 1 + 10 + + + 500MB + + + + + + + + + + + + + \ No newline at end of file diff --git a/clearing-parent/swt-exporter/src/test/java/ru/spcex/clearing/swt/exporter/AbstractServiceTest.java b/clearing-parent/swt-exporter/src/test/java/ru/spcex/clearing/swt/exporter/AbstractServiceTest.java new file mode 100644 index 000000000..8c2fd1d31 --- /dev/null +++ b/clearing-parent/swt-exporter/src/test/java/ru/spcex/clearing/swt/exporter/AbstractServiceTest.java @@ -0,0 +1,53 @@ +package ru.spcex.clearing.swt.exporter; + +import org.apache.kafka.clients.producer.Producer; +import org.junit.jupiter.api.extension.ExtendWith; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.test.context.ContextConfiguration; +import org.springframework.test.context.junit.jupiter.SpringExtension; +import ru.clearing.classes.statics.data.registry.Registry; +import ru.clearing.classes.statics.data.registry.TradingClearingRegistry; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; +import ru.spcex.clearing.test.config.ImdgTestConfig; +import ru.spcex.clearing.test.config.KafkaTestConfig; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; + +import java.time.LocalDate; + +import static ru.spcex.clearing.test.config.ImdgTestConfig.currentID; +import static ru.spcex.clearing.test.config.ImdgTestConfig.waitAvailableImdgProviderAndAddAdminWithDefaultId; + +@ExtendWith(SpringExtension.class) +@ContextConfiguration(classes = { + ImdgTestConfig.class, + KafkaTestConfig.class}) +public abstract class AbstractServiceTest { + protected static final long id = currentID.getAndIncrement(); + protected Imdg registryImdg; + protected Imdg tradingClearingRegistryImdg; + protected LocalDate currentDate = LocalDate.now(); + protected String tcrA = "1324A234"; + protected String tcrD = "124324A234"; + protected Long securityIdFirst = 12L; + protected Long securityIdSecond = 23L; + + @Autowired + @Qualifier("mockProducer") + protected Producer mockProducer; + + @Autowired + protected KafkaSender kafkaSender; + + @Autowired + @Qualifier("hazelcastServiceTest") + protected ImdgProvider imdgProvider; + + protected void init() { + waitAvailableImdgProviderAndAddAdminWithDefaultId(); + this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class); + this.tradingClearingRegistryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class); + } +} diff --git a/clearing-parent/swt-exporter/src/test/java/ru/spcex/clearing/swt/exporter/services/AbstractExporterServiceTest.java b/clearing-parent/swt-exporter/src/test/java/ru/spcex/clearing/swt/exporter/services/AbstractExporterServiceTest.java new file mode 100644 index 000000000..200ca49b8 --- /dev/null +++ b/clearing-parent/swt-exporter/src/test/java/ru/spcex/clearing/swt/exporter/services/AbstractExporterServiceTest.java @@ -0,0 +1,23 @@ +package ru.spcex.clearing.swt.exporter.services; + +import org.junit.jupiter.api.Test; +import ru.spcex.clearing.swt.exporter.AbstractServiceTest; +import ru.spcex.clearing.swt.exporter.services.exportimpl.MoneyExporterService; + +import javax.annotation.PostConstruct; + +class AbstractExporterServiceTest extends AbstractServiceTest { + @PostConstruct + public void init() { + super.init(); + } + + @Test + void sendSwtExportedNotification() { + AbstractExporterService moneyExporterService = new MoneyExporterService(null, kafkaSender, imdgProvider); + + String fileName = "KS_RDC_DF-14_fund_202305241832.swt"; + moneyExporterService.sendSwtxportedNotification(fileName); + //TestUtils.waitingSendAndCheckRecord(null, mockProducer); + } +} \ No newline at end of file diff --git a/clearing-parent/swt-exporter/src/test/java/ru/spcex/clearing/swt/exporter/services/MoneyExporterServiceTest.java b/clearing-parent/swt-exporter/src/test/java/ru/spcex/clearing/swt/exporter/services/MoneyExporterServiceTest.java new file mode 100644 index 000000000..e84a87867 --- /dev/null +++ b/clearing-parent/swt-exporter/src/test/java/ru/spcex/clearing/swt/exporter/services/MoneyExporterServiceTest.java @@ -0,0 +1,102 @@ +package ru.spcex.clearing.swt.exporter.services; + +import org.junit.jupiter.api.Test; +import org.springframework.util.StringUtils; +import ru.clearing.classes.statics.data.registry.Registry; +import ru.clearing.classes.statics.data.registry.TradingClearingRegistry; +import ru.spcex.clearing.swt.exporter.AbstractServiceTest; +import ru.spcex.clearing.swt.exporter.services.exportimpl.MoneyExporterService; + +import javax.annotation.PostConstruct; +import java.math.BigDecimal; +import java.util.Collection; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static ru.spcex.clearing.test.TestUtils.clearAllInImdg; + +class MoneyExporterServiceTest extends AbstractServiceTest { + @PostConstruct + public void init() { + super.init(); + } + + /** + * {@link MoneyExporterService#getSwtFileRows()}
+ * Тест проверяет создание строк документа lim.
+ */ + @Test // todo rewrite + void getSwtFileRows() { + clearAllInImdg(tradingClearingRegistryImdg); + Registry registryA = getRegistryA(tcrA, securityIdFirst); + Registry registryD = getRegistryD(tcrA, securityIdFirst); + registryImdg.insert(registryA); + registryImdg.insert(registryD); + registryA = getRegistryA(tcrD, securityIdSecond); + registryD = getRegistryD(tcrD, securityIdSecond); + registryImdg.insert(registryA); + registryImdg.insert(registryD); + + MoneyExporterService moneyExporterService = new MoneyExporterService(null, kafkaSender, imdgProvider); + Collection limFileRows = moneyExporterService.getLimFileRows(); + assertEquals(0, limFileRows.size()); + + TradingClearingRegistry tradingClearingRegistry = new TradingClearingRegistry(); + tradingClearingRegistry.setId(securityIdFirst); + tradingClearingRegistry.setStatus("ACTV"); + tradingClearingRegistryImdg.insert(tradingClearingRegistry); + tradingClearingRegistry.setId(securityIdSecond); + tradingClearingRegistryImdg.insert(tradingClearingRegistry); + + limFileRows = moneyExporterService.getLimFileRows(); + assertEquals(2, limFileRows.size()); + assertTrue(limFileRows.contains(moneyExporterService.getRow(registryA, registryD))); + } + + private Registry getRegistryA(String tradingClearingRegistry, Long securityId) { + Registry registry = new Registry(); + registry.setTradingCode("1A12323"); + registry.setBalance(new BigDecimal("10.00")); + registry.setTradingClearingRegistry(tradingClearingRegistry); + registry.setRegistryDesignation("A"); + registry.setRegistryInstrumentType("M"); + registry.setRegistryUnit("F"); + registry.setClearingCode(clearingCode(registry)); + registry.setClearingDate(currentDate); + registry.setSecuritySymbol("RUB"); + registry.setSecurityId(securityId); + registry.setTradingClearingRegistryId(securityId); + return registry; + } + + private Registry getRegistryD(String tradingClearingRegistry, Long securityId) { + Registry registry = new Registry(); + registry.setTradingCode("1A12323"); + registry.setBalance(new BigDecimal("5.00")); + registry.setTradingClearingRegistry(tradingClearingRegistry); + registry.setRegistryDesignation("D"); + registry.setRegistryInstrumentType("M"); + registry.setRegistryUnit("T"); + registry.setClearingCode(clearingCode(registry)); + registry.setClearingDate(currentDate); + registry.setSecuritySymbol("RUB"); + registry.setSecurityId(securityId); + registry.setTradingClearingRegistryId(securityId); + return registry; + } + + + // see clearing-service ReistryUtil: + + public static String clearingCode(Registry ofRegistry) { + return clearingCode(ofRegistry.getRegistryDesignation(), ofRegistry.getRegistryInstrumentType(), ofRegistry.getRegistryCapacity(), ofRegistry.getRegistryUnit()); + } + + public static String clearingCode(String registryDesignation, String registryInstrumentType, String registryCapacity, String registryUnit) { + if (StringUtils.isEmpty(registryDesignation)) registryDesignation = "-"; + if (StringUtils.isEmpty(registryInstrumentType)) registryInstrumentType = "-"; + if (StringUtils.isEmpty(registryCapacity)) registryCapacity = "-"; + if (StringUtils.isEmpty(registryUnit)) registryUnit = "-"; + return registryDesignation + registryInstrumentType + registryCapacity + registryUnit; + } +} \ No newline at end of file