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/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..28f95aace --- /dev/null +++ b/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/config/settings/ExportSwtServiceSettings.java @@ -0,0 +1,53 @@ +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 String docOut; + //todo ??? private Long interval + + + 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 KafkaProducerSettings getKafkaProducer() { + return kafkaProducer; + } + + public void setKafkaProducer(KafkaProducerSettings kafkaProducer) { + this.kafkaProducer = kafkaProducer; + } + + public String getDocOut() { + return docOut; + } + + public void setDocOut(String docOut) { + this.docOut = docOut; + } +} 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..aac068e27 --- /dev/null +++ b/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/services/AbstractExporterService.java @@ -0,0 +1,156 @@ +package ru.spcex.clearing.swt.exporter.services; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import ru.spcex.clearing.platform.messaging.domain.cud.utilities.JournalEventExportedRequest; +import ru.spcex.clearing.platform.messaging.serialization.LogFormatter; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; +import ru.spcex.platform.classes.base.SpcexObjectBase; +import ru.spcex.platform.enumeration.ResultStatuses; +import ru.spcex.platform.enumeration.SwtTable; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; + +import java.io.*; +import java.time.LocalDateTime; +import java.time.format.DateTimeFormatter; +import java.util.*; + +import static ru.spcex.clearing.platform.messaging.domain.Consts.JOURNAL_SERVICE; + +public abstract class AbstractExporterService { + protected final Logger log = LoggerFactory.getLogger(getClass()); + protected final SwtTable type; + protected final Imdg sdfImdg; + private final DateTimeFormatter dtFileNameFormatter = DateTimeFormatter.ofPattern("yyyyMMddHHmmss"); + private final DateTimeFormatter dtInFileHeaderFormatter = DateTimeFormatter.ofPattern("yyyyMMdd'/'HHmm"); + private final KafkaSender kafkaSender; + protected final FileStorage fileStorage; + + protected AbstractExporterService(FileStorage fileStorage, + KafkaSender kafkaSender, ImdgProvider imdgProvider, + SwtTable type, + String mapName, Class mapClass) { + this.fileStorage = fileStorage; + this.kafkaSender = kafkaSender; + this.type = type; + Objects.requireNonNull(type, "SWT table type not set"); + this.sdfImdg = imdgProvider.getImdg(mapName, mapClass); + } + + public SwtTable getType() { + return type; + } + + protected abstract String typeForFileName(); + + protected abstract String sectionForFileName(); + + protected abstract String typeForHeader(); + + public void process() { + LocalDateTime exportAt = LocalDateTime.now(); + + String fileName = formatFileName(typeForFileName(), sectionForFileName(), exportAt); + log.debug("Start export {} Lim file", fileName); + + try { + byte[] data; + { + ByteArrayOutputStream outBuffer = new ByteArrayOutputStream(); + Collection records = selectItems(); + log.debug("Prepared {} record from {} to file {}", + records.size(), sdfImdg.getMapName(), fileName); + makeSWTData(typeForHeader(), exportAt, records, outBuffer); + data = outBuffer.toByteArray(); + } + fileStorage.saveFile(fileName, data); + } catch (Exception e) { // IOException, ... + log.error("Failed export {} file", fileName); + sendSwtExportedNotification(exportAt, null, ResultStatuses.notSuccess); + throw new RuntimeException(e); + } + log.debug("Successfully exported {} file", fileName); + + sendSwtExportedNotification(exportAt, null, ResultStatuses.success); + } + + protected abstract String getDocumentNameForJournal(); + + void sendSwtExportedNotification(LocalDateTime registrationAt, Long registrationNumber, ResultStatuses resultStatus) { + JournalEventExportedRequest exportedRequest = new JournalEventExportedRequest(); + exportedRequest.setRegistratoinDate(registrationAt.toLocalDate()); + exportedRequest.setRegistrationTime(registrationAt.toLocalTime()); + exportedRequest.setRegistrationNumber(registrationNumber); + exportedRequest.setDocumentName(getDocumentNameForJournal()); + exportedRequest.setResultStatus(resultStatus.getKey()); + log.debug("Send message to kafka \"{}\": {}", JOURNAL_SERVICE, LogFormatter.toStringWrapper(exportedRequest)); + kafkaSender.sendRequestToQueue(JOURNAL_SERVICE, exportedRequest); + } + + + /** + * @param type DF-09 + * @param section bond/fund/"" + * @param atTime LocalDateTime.now(), если null - текущее время + * @return пример "KS_RDC_DF-12_bond_220907151804503.txt" + */ + protected String formatFileName(String type, String section, LocalDateTime atTime) { + Objects.requireNonNull(type); + if (section == null) section = ""; + if (section.length() > 0) section += "_"; + if (atTime == null) atTime = LocalDateTime.now(); + String dt = dtFileNameFormatter.format(atTime); + return String.format("KS_RDC_%s_%s%s.txt", type, section, dt); + } + + // Выборка + protected Collection selectItems() { + return sdfImdg.getAllValues(); + } + + // Конвертация (поля см. meta.xml) + protected abstract LinkedHashMap convertRecord(T record); + +// protected abstract String[] swtHeader(); + + protected void makeSWTData(String type, LocalDateTime time, Collection records, OutputStream outStream) { + PrintWriter out = new PrintWriter(outStream); + // todo SWT txt не понял формат. Надо уточнить формат файла. Должен соответствовать мете. + out.println("To:CSO"); + out.println("From:SPCE"); + if (type != null) + out.println("Type:" + type);// Type:009 + if (time != null) { + String timeS = dtInFileHeaderFormatter.format(time); + out.println("Date/Time:" + timeS);// Date/Time:20230227/0932 + } + /* + To:CSO + From:SPCE + Type:009 + Date/Time:20230227/0932 + :20:0ef63e17-83ac-4e22-a76b-fd8a4df10de3 + :21:SDC230227084839 + :18A:46 +*/ + StringBuilder line = new StringBuilder(); + for (T row : records) { + LinkedHashMap rowData = convertRecord(row); + line.setLength(0); + for (Map.Entry r : rowData.entrySet()) { + String value = convertItem(r.getValue()); + line.append(value).append(':'); + } + if (line.length() > 0) // remove : + line.setLength(line.length() - 1); + out.println(line); + } + //out.println("2"); // todo что значит 2? + } + + protected String convertItem(Object o) { + //todo date/time/etc. + return String.valueOf(o); + } +} 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..fe29dd157 --- /dev/null +++ b/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/services/FileStorage.java @@ -0,0 +1,37 @@ +package ru.spcex.clearing.swt.exporter.services; + +import org.apache.commons.io.FileUtils; +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.getDocOut() == null || config.getDocOut().isBlank()) { + throw new IllegalArgumentException("Out directory settings is empty."); + } + this.outPath = new File(config.getDocOut()); + 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 { + File toFile = new File(outPath, fileName); + FileUtils.writeByteArrayToFile(toFile, data); + } +} 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..ec1b92f09 --- /dev/null +++ b/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/services/LauncherCommandReceiver.java @@ -0,0 +1,70 @@ +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.BaseRequest; +import ru.spcex.clearing.platform.messaging.domain.Consts; +import ru.spcex.clearing.platform.messaging.domain.cud.importexport.SwtExporterRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest; +import ru.spcex.clearing.platform.messaging.service.QueueConsumer; +import ru.spcex.platform.utils.log.ExceptionUtils; + +import java.util.List; + +@Service +public class LauncherCommandReceiver extends QueueConsumer implements InitializingBean { + private final Logger log = LoggerFactory.getLogger(getClass()); + protected final List exporterServices; + + public LauncherCommandReceiver(Consumer kafkaQueue, + List exporterServices + ) { + super(kafkaQueue); + this.exporterServices = exporterServices; + } + + @Override + public void afterPropertiesSet() { +// callback(LauncherCommandRequest.class) +// .setConsumer(this::exportAll) +// .forDestination(Task.unloadingSession_LIMM.topic(), callbacks::put); // todo task name? + callback(SwtExporterRequest.class) + .setConsumer(this::exportSpecial) + .forDestination(Consts.SWT_EXPORTER, callbacks::put); + init(); + } + + protected void exportAll(BaseRequest request) { + log.info("LauncherCommandRequest request received: {}", request); + for (AbstractExporterService exporter : exporterServices) { + log.debug("Export {}", exporter); + try { + exporter.process(); + } catch (Exception e) { + log.error("One of exporter has error: {}", ExceptionUtils.getStackTrace(e)); + } + } + log.info("All SWT export has finished."); + } + + protected void exportSpecial(BaseRequest request) { + log.info("SwtExporterRequest request received: {}", request); + SwtExporterRequest req = request.getRequestPayload(); + for (AbstractExporterService exporter : exporterServices) { + if (exporter.getType()==req.getType()) { + log.debug("Export {}", exporter); + try { + exporter.process(); + } catch (Exception e) { + log.error("Exporter has error: {}", ExceptionUtils.getStackTrace(e)); + } + } + return; + } + log.error("SWT export not execute for type \"{}\" - unknown command", req.getType()); + //todo return error? + } +} diff --git a/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/services/exportimpl/DF09Exporter.java b/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/services/exportimpl/DF09Exporter.java new file mode 100644 index 000000000..8334eefef --- /dev/null +++ b/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/services/exportimpl/DF09Exporter.java @@ -0,0 +1,62 @@ +package ru.spcex.clearing.swt.exporter.services.exportimpl; + +import org.springframework.stereotype.Service; +import ru.clearing.classes.statics.data.sdf.SDf09; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +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.enumeration.SwtTable; +import ru.spcex.platform.imdg.api.ImdgProvider; + +import java.util.LinkedHashMap; + +@Service +public class DF09Exporter extends AbstractExporterService { + + public DF09Exporter(FileStorage fileStorage, + KafkaSender kafkaSender, + ImdgProvider imdgProvider) { + super(fileStorage, kafkaSender, imdgProvider, + SwtTable.SDF_09, + IMDGDistributedNames.Map_SDf09, SDf09.class); + } + + @Override + protected String typeForFileName() { + return "DF-09"; + } + + @Override + protected String sectionForFileName() { + return null; + } + + @Override + protected String typeForHeader() { + return "009"; + } + + @Override + protected String getDocumentNameForJournal() { + return "Уведомление об исполнении операции загрузки ценных бумаг или уведомление об ошибке"; + } + + @Override + protected LinkedHashMap convertRecord(SDf09 record) { + LinkedHashMap row = new LinkedHashMap<>(); + row.put("id", record.getId()); + + row.put("outDocument", record.getOutDocument()); + row.put("inDocument", record.getInDocument()); + row.put("depoCode", record.getDepoCode()); + row.put("quantity", record.getQuantity()); + row.put("securityCode", record.getSecurityCode()); + row.put("clientName", record.getClientName()); + row.put("result", record.getResult()); + row.put("generationTime", record.getGenerationTime()); + row.put("generationId", record.getGenerationId()); + return row; + } + +} diff --git a/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/services/exportimpl/DF11Exporter.java b/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/services/exportimpl/DF11Exporter.java new file mode 100644 index 000000000..53581fe6e --- /dev/null +++ b/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/services/exportimpl/DF11Exporter.java @@ -0,0 +1,63 @@ +package ru.spcex.clearing.swt.exporter.services.exportimpl; + +import org.springframework.stereotype.Service; +import ru.clearing.classes.statics.data.sdf.SDf11; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +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.enumeration.SwtTable; +import ru.spcex.platform.imdg.api.ImdgProvider; + +import java.util.LinkedHashMap; + + +@Service +public class DF11Exporter extends AbstractExporterService { + + public DF11Exporter(FileStorage fileStorage, + KafkaSender kafkaSender, + ImdgProvider imdgProvider) { + super(fileStorage, kafkaSender, imdgProvider, + SwtTable.SDF_11, + IMDGDistributedNames.Map_SDf11, SDf11.class); + } + + @Override + protected String typeForFileName() { + return "DF-11"; + } + + @Override + protected String sectionForFileName() { + return null; + } + + @Override + protected String typeForHeader() { + return "011"; + } + + @Override + protected String getDocumentNameForJournal() { + return "Ответ на Запрос на Зачисление или списание ценных бумаг"; + } + + @Override + protected LinkedHashMap convertRecord(SDf11 record) { + LinkedHashMap row = new LinkedHashMap<>(); + row.put("id", record.getId()); + + row.put("outDocument", record.getOutDocument()); + row.put("inDocument", record.getInDocument()); + row.put("depoCode", record.getDepoCode()); + row.put("quantity", record.getQuantity()); + row.put("securityCode", record.getSecurityCode()); + row.put("clientName", record.getClientName()); + row.put("result", record.getResult()); + row.put("generationTime", record.getGenerationTime()); + row.put("generationId", record.getGenerationId()); + return row; + } + +} diff --git a/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/services/exportimpl/DF12Exporter.java b/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/services/exportimpl/DF12Exporter.java new file mode 100644 index 000000000..7892880e7 --- /dev/null +++ b/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/services/exportimpl/DF12Exporter.java @@ -0,0 +1,66 @@ +package ru.spcex.clearing.swt.exporter.services.exportimpl; + +import org.springframework.stereotype.Service; +import ru.clearing.classes.statics.data.sdf.SDf12; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +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.enumeration.SwtTable; +import ru.spcex.platform.imdg.api.ImdgProvider; + +import java.util.LinkedHashMap; + + +@Service +public class DF12Exporter extends AbstractExporterService { + + public DF12Exporter(FileStorage fileStorage, + KafkaSender kafkaSender, + ImdgProvider imdgProvider) { + super(fileStorage, kafkaSender, imdgProvider, + SwtTable.SDF_12, + IMDGDistributedNames.Map_SDf12, SDf12.class); + } + + @Override + protected String typeForFileName() { + return "DF-12"; + } + + @Override + protected String sectionForFileName() { + return "bond"; // todo bond / fund + } + + @Override + protected String typeForHeader() { + return "012"; + } + + @Override + protected String getDocumentNameForJournal() { + //todo Выбор: + // Распоряжение на проведение операций по итогам клиринга (Фондовая секция) + //или + // Распоряжение на проведение операций по итогам клиринга (ОФЗ, ОБР) + return "Распоряжение на проведение операций по итогам клиринга (Фондовая секция / ОФЗ, ОБР)"; + } + + @Override + protected LinkedHashMap convertRecord(SDf12 record) { + LinkedHashMap row = new LinkedHashMap<>(); + row.put("id", record.getId()); + + row.put("outDocument", record.getOutDocument()); + row.put("quantity", record.getQuantity()); + row.put("securityCode", record.getSecurityCode()); + row.put("depoCodeSender", record.getDepoCodeSender()); + row.put("depoCodeAdressee", record.getDepoCodeAdressee()); + row.put("result", record.getResult()); + row.put("generationTime", record.getGenerationTime()); + row.put("generationId", record.getGenerationId()); + return row; + } + +} diff --git a/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/services/exportimpl/DF14Exporter.java b/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/services/exportimpl/DF14Exporter.java new file mode 100644 index 000000000..c2a706a61 --- /dev/null +++ b/clearing-parent/swt-exporter/src/main/java/ru/spcex/clearing/swt/exporter/services/exportimpl/DF14Exporter.java @@ -0,0 +1,62 @@ +package ru.spcex.clearing.swt.exporter.services.exportimpl; + +import org.springframework.stereotype.Service; +import ru.clearing.classes.statics.data.sdf.SDf14; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +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.enumeration.SwtTable; +import ru.spcex.platform.imdg.api.ImdgProvider; + +import java.util.LinkedHashMap; + + +@Service +public class DF14Exporter extends AbstractExporterService { + + public DF14Exporter(FileStorage fileStorage, + KafkaSender kafkaSender, + ImdgProvider imdgProvider) { + super(fileStorage, kafkaSender, imdgProvider, + SwtTable.SDF_14, + IMDGDistributedNames.Map_SDf14, SDf14.class); + } + + @Override + protected String typeForFileName() { + return "DF-14"; + } + + @Override + protected String sectionForFileName() { + return "bond"; // todo bond / fund + } + + @Override + protected String typeForHeader() { + return "014"; + } + + @Override + protected String getDocumentNameForJournal() { + //todo Выбор: + // Уведомление о завершении расчетов (ОФЗ, ОБР) + //или + // Уведомление о завершении расчетов (ОФЗ, ОБР) + return "Уведомление о завершении расчетов (ОФЗ, ОБР)"; + } + + @Override + protected LinkedHashMap convertRecord(SDf14 record) { + LinkedHashMap row = new LinkedHashMap<>(); + row.put("id", record.getId()); + + row.put("outDocument", record.getOutDocument()); + row.put("result", record.getResult()); + row.put("generationTime", record.getGenerationTime()); + row.put("generationId", record.getGenerationId()); + return row; + } + +} 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..263f91c96 --- /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.docOut=/opt/spcex/clearing/files/swt/SettlementHouse_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..db500402c --- /dev/null +++ b/clearing-parent/swt-exporter/src/test/java/ru/spcex/clearing/swt/exporter/services/AbstractExporterServiceTest.java @@ -0,0 +1,25 @@ +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.DF09Exporter; +import ru.spcex.platform.enumeration.ResultStatuses; + +import javax.annotation.PostConstruct; +import java.time.LocalDateTime; + +class AbstractExporterServiceTest extends AbstractServiceTest { + @PostConstruct + public void init() { + super.init(); + } + + @Test + void sendSwtExportedNotification() { + AbstractExporterService moneyExporterService = new DF09Exporter(null, kafkaSender, imdgProvider); + + String fileName = "KS_RDC_DF-14_bond_221005134616035.txt"; + moneyExporterService.sendSwtExportedNotification(LocalDateTime.now(), 1L, ResultStatuses.notSuccess); + //TestUtils.waitingSendAndCheckRecord(null, mockProducer); + } +} \ No newline at end of file diff --git a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/SwtTable.java b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/SwtTable.java new file mode 100644 index 000000000..f650cda65 --- /dev/null +++ b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/SwtTable.java @@ -0,0 +1,24 @@ +package ru.spcex.platform.enumeration; + +import ru.spcex.platform.utils.enumeration.IEnumKey; + +public enum SwtTable implements IEnumKey { + SDF_09("SDF_09"), SDF_11("SDF_11"), + SDF_12("SDF_12"), SDF_14("SDF_14"); + + SwtTable(String key) { + this.key = key; + } + + private final String key; + + @Override + public String getKey() { + return this.key; + } + + @Override + public boolean equalsByKey(String key) { + return IEnumKey.super.equalsByKey(key); + } +} diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java index 2c8a7c4aa..00a27f108 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java @@ -123,11 +123,13 @@ public interface Consts { String SDF56_PROCESS = "sdf56-process"; String CONTINUE_SESSION_BN_FIRST_PART = "sdf57-process"; String CONTINUE_SESSION_BN_SECOND_PART = "sdf13-process"; + String SWT_EXPORTER = "swt-exporter"; String REVISE_PROCESS = "revise-process"; String EXPORT_PROCESS = "export-process"; String EXPORT_COMPLETED = "export_completed"; String S_TRADES_IMPORTED = "s_trades-imported"; String LIM_EXPORTED = "lim_exported"; + String JOURNAL_SERVICE = "journal-service-exported"; String ACCOUNT_TERMINATION = "account-termination"; String BALANCE_ACCOUNT_NEW = "balance-account-new"; String BALANCE_ACCOUNT_UPDATE = "balance-account-update"; diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/importexport/SwtExporterRequest.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/importexport/SwtExporterRequest.java new file mode 100644 index 000000000..41b54f2e5 --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/importexport/SwtExporterRequest.java @@ -0,0 +1,17 @@ +package ru.spcex.clearing.platform.messaging.domain.cud.importexport; + +import com.fasterxml.jackson.annotation.JsonProperty; +import ru.spcex.platform.enumeration.SwtTable; + +public class SwtExporterRequest { + @JsonProperty + public SwtTable type; + + public SwtTable getType() { + return type; + } + + public void setType(SwtTable type) { + this.type = type; + } +} diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/utilities/JournalEventExportedRequest.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/utilities/JournalEventExportedRequest.java new file mode 100644 index 000000000..099c346a4 --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/utilities/JournalEventExportedRequest.java @@ -0,0 +1,83 @@ +package ru.spcex.clearing.platform.messaging.domain.cud.utilities; + +import com.fasterxml.jackson.annotation.JsonProperty; +import com.fasterxml.jackson.databind.annotation.JsonDeserialize; +import com.fasterxml.jackson.databind.annotation.JsonSerialize; +import ru.spcex.clearing.platform.messaging.domain.json.deserialize.LocalDateDeserializer; +import ru.spcex.clearing.platform.messaging.domain.json.deserialize.LocalTimeDeserializer; +import ru.spcex.clearing.platform.messaging.domain.json.serialize.LocalDateSerializer; +import ru.spcex.clearing.platform.messaging.domain.json.serialize.LocalTimeSerializer; + +import java.time.LocalDate; +import java.time.LocalTime; + +public class JournalEventExportedRequest { + @JsonSerialize(using = LocalDateSerializer.class) + @JsonDeserialize(using = LocalDateDeserializer.class) + @JsonProperty + private LocalDate registratoinDate; + @JsonSerialize(using = LocalTimeSerializer.class) + @JsonDeserialize(using = LocalTimeDeserializer.class) + @JsonProperty + private LocalTime registrationTime; + @JsonProperty + private Long registrationNumber; + @JsonProperty + private String documentName; + @JsonProperty + private String dossierNumber; + /** + * ACK при успешной загрузке + * NACK при ошибке загрузки + */ + @JsonProperty + private String resultStatus; + + public LocalDate getRegistratoinDate() { + return registratoinDate; + } + + public void setRegistratoinDate(LocalDate registratoinDate) { + this.registratoinDate = registratoinDate; + } + + public LocalTime getRegistrationTime() { + return registrationTime; + } + + public void setRegistrationTime(LocalTime registrationTime) { + this.registrationTime = registrationTime; + } + + public Long getRegistrationNumber() { + return registrationNumber; + } + + public void setRegistrationNumber(Long registrationNumber) { + this.registrationNumber = registrationNumber; + } + + public String getDocumentName() { + return documentName; + } + + public void setDocumentName(String documentName) { + this.documentName = documentName; + } + + public String getDossierNumber() { + return dossierNumber; + } + + public void setDossierNumber(String dossierNumber) { + this.dossierNumber = dossierNumber; + } + + public String getResultStatus() { + return resultStatus; + } + + public void setResultStatus(String resultStatus) { + this.resultStatus = resultStatus; + } +} diff --git a/pom.xml b/pom.xml index ce3a0c7c7..6db1b7fb2 100644 --- a/pom.xml +++ b/pom.xml @@ -37,6 +37,8 @@ ${folder_root_clearing}/clearing-parent/imdg ${folder_root_clearing}/clearing-parent/dbf-exporter ${folder_root_clearing}/clearing-parent/dbf-importer + ${folder_root_clearing}/clearing-parent/lim-exporter + ${folder_root_clearing}/clearing-parent/swt-exporter ${folder_root_clearing}/clearing-parent/trade-importer ${folder_root_clearing}/clearing-parent/account-service ${folder_root_clearing}/clearing-parent/balance-service diff --git a/z-distr/pom.xml b/z-distr/pom.xml index 819183670..065065934 100644 --- a/z-distr/pom.xml +++ b/z-distr/pom.xml @@ -191,6 +191,45 @@ + + + copy-lim-exporter-bin + prepare-package + + copy + + + + + ${folder_root_lim-exporter}/target/lim-exporter.jar + ${folder.clearing.distr.modules}/lim-exporter/lim-exporter.jar + + + ${folder_root_lim-exporter}/src/main/resources/application.properties + ${folder.clearing.distr.modules}/lim-exporter/application.properties + + + + + + copy-swt-exporter-bin + prepare-package + + copy + + + + + ${folder_root_swt-exporter}/target/swt-exporter.jar + ${folder.clearing.distr.modules}/swt-exporter/swt-exporter.jar + + + ${folder_root_swt-exporter}/src/main/resources/application.properties + ${folder.clearing.distr.modules}/swt-exporter/application.properties + + + + copy-trade-importer-bin prepare-package diff --git a/z-distr/src/main/resources/distr/bin/kill_all.sh b/z-distr/src/main/resources/distr/bin/kill_all.sh index ee3ee0994..05a6100e4 100644 --- a/z-distr/src/main/resources/distr/bin/kill_all.sh +++ b/z-distr/src/main/resources/distr/bin/kill_all.sh @@ -7,6 +7,7 @@ kill -9 $(ps -ef | grep java | grep company-service.jar | awk '{print $2}') kill -9 $(ps -ef | grep java | grep clearing-service.jar | awk '{print $2}') kill -9 $(ps -ef | grep java | grep dbf-exporter.jar | awk '{print $2}') kill -9 $(ps -ef | grep java | grep lim-exporter.jar | awk '{print $2}') +kill -9 $(ps -ef | grep java | grep swt-exporter.jar | awk '{print $2}') kill -9 $(ps -ef | grep java | grep dbf-importer.jar | awk '{print $2}') kill -9 $(ps -ef | grep java | grep trade-importer.jar | awk '{print $2}') kill -9 $(ps -ef | grep java | grep imdg.jar | awk '{print $2}') diff --git a/z-distr/src/main/resources/distr/bin/swt-exporter.sh b/z-distr/src/main/resources/distr/bin/swt-exporter.sh new file mode 100644 index 000000000..1f22b4caa --- /dev/null +++ b/z-distr/src/main/resources/distr/bin/swt-exporter.sh @@ -0,0 +1,9 @@ +#!/bin/bash + +CLEARING_HOME=/opt/mfd/clearing/ +cd $CLEARING_HOME/bin + +CMD="java -Xrunjdwp:transport=dt_socket,server=y,suspend=n,address=7120 -jar swt-exporter.jar --spring.config.location=$CLEARING_HOME/settings/swt-exporter/" + +$CMD >/dev/null 2>&1 & +