From ef9f9b3ca2d440d19cd5739c7a186b9367bdd135 Mon Sep 17 00:00:00 2001 From: psemenkov Date: Mon, 6 May 2024 10:55:52 +0300 Subject: [PATCH] =?UTF-8?q?http://jira.mfd.msk:8088/browse/CLS-665=20?= =?UTF-8?q?=D0=94=D0=BE=D0=B1=D0=B0=D0=B2=D0=B8=D0=BB=20=D0=BB=D0=BE=D0=B3?= =?UTF-8?q?=D0=B8=D1=80=D0=BE=D0=B2=D0=B0=D0=BD=D0=B8=D0=B5=20=D0=BF=D0=B0?= =?UTF-8?q?=D1=80=20sdf=20=D0=B2=20pairsdf.?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../logic/stages/ExportFromHazelcast.java | 25 +++- .../dbf/importer/logic/data/enums/ETable.java | 4 + .../logic/data/tables/AbstractTable.java | 8 ++ .../logic/stages/DbfImportKafkaMessenger.java | 13 ++ .../dbf/importer/logic/stages/ImportToDB.java | 5 + clearing-parent/utility-service/pom.xml | 5 + .../ru/spcex/clearing/utility/Enum/Table.java | 75 ++++++++++++ .../utility/service/PairSdfService.java | 113 ++++++++++++++++++ .../HazelcastInstanceTestConfiguration.java | 23 ---- .../HazelcastServiceTestConfiguration.java | 73 ----------- .../clearing/utility/config/TestConfig.java | 16 +++ .../service/KeyRateServiceDisabled.java | 9 -- .../utility/service/PairSdfServiceTest.java | 103 ++++++++++++++++ .../platform/messaging/domain/Consts.java | 1 + .../domain/cud/utilities/PairSdfRequest.java | 36 ++++++ 15 files changed, 400 insertions(+), 109 deletions(-) create mode 100644 clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/Enum/Table.java create mode 100644 clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/service/PairSdfService.java delete mode 100644 clearing-parent/utility-service/src/test/java/ru/spcex/clearing/utility/config/HazelcastInstanceTestConfiguration.java delete mode 100644 clearing-parent/utility-service/src/test/java/ru/spcex/clearing/utility/config/HazelcastServiceTestConfiguration.java create mode 100644 clearing-parent/utility-service/src/test/java/ru/spcex/clearing/utility/config/TestConfig.java create mode 100644 clearing-parent/utility-service/src/test/java/ru/spcex/clearing/utility/service/PairSdfServiceTest.java create mode 100644 platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/utilities/PairSdfRequest.java diff --git a/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/logic/stages/ExportFromHazelcast.java b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/logic/stages/ExportFromHazelcast.java index d94032686..482e4aef8 100644 --- a/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/logic/stages/ExportFromHazelcast.java +++ b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/logic/stages/ExportFromHazelcast.java @@ -1,6 +1,5 @@ package ru.spcex.clearing.dbf.exporter.logic.stages; -import com.linuxense.javadbf.DBFField; import com.linuxense.javadbf.DBFWriter; import org.springframework.beans.factory.InitializingBean; import org.springframework.stereotype.Component; @@ -11,6 +10,9 @@ import ru.spcex.clearing.dbf.exporter.logic.data.ResultContainer; import ru.spcex.clearing.dbf.exporter.logic.data.enums.StageResult; import ru.spcex.clearing.dbf.exporter.logic.data.enums.Table; import ru.spcex.clearing.dbf.exporter.services.converters.DFConverter; +import ru.spcex.clearing.platform.messaging.domain.cud.utilities.PairSdfRequest; +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.CurrencyCode; import ru.spcex.platform.imdg.api.Imdg; @@ -19,8 +21,10 @@ import ru.spcex.platform.imdg.api.ImdgProvider; import java.io.File; import java.nio.charset.Charset; import java.util.*; +import java.util.function.Supplier; import static ru.spcex.clearing.dbf.exporter.services.converters.DFConverter.ruCurrency; +import static ru.spcex.clearing.platform.messaging.domain.Consts.PAIR_SDF; /** * Выгрузка данных из мапы hazelcast и их запись в файлы @@ -31,19 +35,20 @@ public class ExportFromHazelcast extends Stage implements InitializingBean { private final ImdgProvider imdgProvider; private final SFTPConfig.DbfGateway gateway; private final List converters; + private final Supplier kafkaSender; - - private final Map dbfFieldsForTable = new HashMap<>(); private Charset dbfCharset; public ExportFromHazelcast(ExportDBFServiceSettings settings, ImdgProvider imdgProvider, SFTPConfig.DbfGateway gateway, - List converters) { + List converters, + Supplier kafkaSender) { this.settings = settings; this.imdgProvider = imdgProvider; this.gateway = gateway; this.converters = converters; + this.kafkaSender = kafkaSender; } @Override @@ -110,6 +115,9 @@ public class ExportFromHazelcast extends Stage implements InitializingBean { } else { log.debug("uuid {}, send to SFTP path \"{}\"", resultContainer.getUuid(), path); gateway.sendToSftp(dbfFile, path); + if (Table.S_DF02.equals(resultContainer.getTableForExport())) { + sendPairSdfRequest(resultContainer); + } } } @@ -133,4 +141,13 @@ public class ExportFromHazelcast extends Stage implements InitializingBean { } } + private void sendPairSdfRequest(ResultContainer resultContainer) { + PairSdfRequest request = new PairSdfRequest(); + request.setGenerationId(resultContainer.getGroupId()); + request.setFileNameSDf(resultContainer.getFileForExport().getName()); + request.setTableSDf(resultContainer.getTableForExport().getFilePrefix()); + Long msgId = kafkaSender.get().sendRequestToQueue(PAIR_SDF, request); + log.info("Send PairSdfRequest={} message id={} to kafka \"{}\"", LogFormatter.toString(request), msgId, PAIR_SDF); + } + } diff --git a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/data/enums/ETable.java b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/data/enums/ETable.java index 0ab232ada..db12e66a2 100644 --- a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/data/enums/ETable.java +++ b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/data/enums/ETable.java @@ -11,6 +11,10 @@ public enum ETable { DF_55("DF-55"), DF_57("DF-57"); + public String getPrefix() { + return prefix; + } + private final String prefix; ETable(String prefix) { diff --git a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/data/tables/AbstractTable.java b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/data/tables/AbstractTable.java index 46e566c02..b5644ab41 100644 --- a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/data/tables/AbstractTable.java +++ b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/data/tables/AbstractTable.java @@ -96,4 +96,12 @@ public abstract class AbstractTable { } throw new IllegalArgumentException("Can not convert to BigDecimal " + value.getClass().getName() + " " + value); } + + public String getFilename() { + return filename; + } + + public Long getFileId() { + return fileId; + } } diff --git a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/stages/DbfImportKafkaMessenger.java b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/stages/DbfImportKafkaMessenger.java index 1fcd0f8b6..3b6bf3287 100644 --- a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/stages/DbfImportKafkaMessenger.java +++ b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/stages/DbfImportKafkaMessenger.java @@ -7,10 +7,12 @@ import org.springframework.stereotype.Component; import ru.spcex.clearing.dbf.importer.logic.data.ResultContainer; import ru.spcex.clearing.dbf.importer.logic.data.enums.ETable; import ru.spcex.clearing.dbf.importer.logic.data.enums.StageResult; +import ru.spcex.clearing.dbf.importer.logic.data.tables.AbstractTable; import ru.spcex.clearing.platform.messaging.domain.Consts; import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest; import ru.spcex.clearing.platform.messaging.domain.cud.clearing.Sdf04Request; import ru.spcex.clearing.platform.messaging.domain.cud.utilities.NotificationNewRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.utilities.PairSdfRequest; import ru.spcex.clearing.platform.messaging.serialization.LogFormatter; import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.platform.enumeration.ObjectType; @@ -22,6 +24,8 @@ import java.util.Map; import java.util.function.Consumer; import java.util.function.Supplier; +import static ru.spcex.clearing.platform.messaging.domain.Consts.PAIR_SDF; + @Component public class DbfImportKafkaMessenger implements InitializingBean { final Logger log = LoggerFactory.getLogger(getClass()); @@ -103,6 +107,15 @@ public class DbfImportKafkaMessenger implements InitializingBean { groupId, table, msgId, destination); } + public void sendPairSdfRequest(AbstractTable table, String tableSdf) { + PairSdfRequest request = new PairSdfRequest(); + request.setGenerationId(table.getFileId()); + request.setFileNameSDf(table.getFilename()); + request.setTableSDf(tableSdf); + Long msgId = kafka.get().sendRequestToQueue(PAIR_SDF, request); + log.info("Send PairSdfRequest={} message id={} to kafka \"{}\"", LogFormatter.toString(request), msgId, PAIR_SDF); + } + private void messageDf04(Long groupId) { Sdf04Request sdf04ImportNotification = new Sdf04Request(); sdf04ImportNotification.setGroupId(groupId); diff --git a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/stages/ImportToDB.java b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/stages/ImportToDB.java index 10a94a7b6..c92470d1b 100644 --- a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/stages/ImportToDB.java +++ b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/stages/ImportToDB.java @@ -17,6 +17,8 @@ import java.io.InputStream; import java.nio.charset.Charset; import java.util.Map; +import static ru.spcex.clearing.dbf.importer.logic.data.enums.ETable.DF_01; + /** * Заливка проверенных данных в базу */ @@ -58,6 +60,9 @@ public class ImportToDB extends Stage { } table.injectEntity(table.getEntity(entity)); } + if (DF_01.equals(currTable)){ + kafkaMessenger.sendPairSdfRequest(table, currTable.getPrefix()); + } kafkaMessenger.notifySystemIfNeeded(currTable, fileId); kafkaMessenger.notifyUserAboutSuccessLoad(resultContainer); } catch (IOException exception) { diff --git a/clearing-parent/utility-service/pom.xml b/clearing-parent/utility-service/pom.xml index 985adde41..66c6e6a55 100644 --- a/clearing-parent/utility-service/pom.xml +++ b/clearing-parent/utility-service/pom.xml @@ -54,6 +54,11 @@ spring-boot-starter-test test + + ru.spcex.clearing + test-clearing + test + diff --git a/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/Enum/Table.java b/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/Enum/Table.java new file mode 100644 index 000000000..b7cb1c11a --- /dev/null +++ b/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/Enum/Table.java @@ -0,0 +1,75 @@ +package ru.spcex.clearing.utility.Enum; + +import ru.clearing.classes.statics.data.sdf.*; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.platform.classes.base.SpcexObjectBase; + +import java.util.Map; +import java.util.function.Function; + +public enum Table { + //from dbf-importer + S_DF01("DF-01", IMDGDistributedNames.Map_SDf01, SDf01.class), + S_DF04("DF-04", IMDGDistributedNames.Map_SDf04, SDf04.class), + S_DF06("DF-06", IMDGDistributedNames.Map_SDf06, SDf06.class), + S_DF52("DF-52", IMDGDistributedNames.Map_SDf52, SDf52.class), + S_DF55("DF-55", IMDGDistributedNames.Map_SDf55, SDf55.class), + //from dbf-exporter + S_DF02("DF-02", IMDGDistributedNames.Map_SDf02, SDf02.class), + S_DF03("DF-03", IMDGDistributedNames.Map_SDf03, SDf03.class), + S_DF07("DF-07", IMDGDistributedNames.Map_SDf07, SDf07.class), + S_DF53("DF-53", IMDGDistributedNames.Map_SDf53, SDf53.class), + S_DF54("DF-54", IMDGDistributedNames.Map_SDf54, SDf54.class); + + public static final Map outTables = Map.of( + S_DF02, new TableProcessor(S_DF01, value -> ((SDf01) value).getGenerationId(), value -> ((SDf02) value).getInSDfId()), + S_DF04, new TableProcessor(S_DF03, value -> ((SDf03) value).getGenerationId(), value -> null), + S_DF07, new TableProcessor(S_DF06, value -> ((SDf06) value).getGenerationId(), value -> ((SDf07) value).getInSDfId()), + S_DF53, new TableProcessor(S_DF52, value -> ((SDf52) value).getGenerationId(), value -> ((SDf53) value).getInSDfId()), + S_DF55, new TableProcessor(S_DF54, value -> ((SDf54) value).getGenerationId(), value -> null)); + + /** + * Имя таблицы + */ + public final String tableSDf; + + /** + * Имя таблицы hazelcast в которой будет поиск generationId + */ + private final String hazelcastMapName; + + /** + * Класс объекта + */ + private final Class entityClass; + + Table(String tableSDf, String hazelcastMapName, Class entityClass) { + this.tableSDf = tableSDf; + this.hazelcastMapName = hazelcastMapName; + this.entityClass = entityClass; + } + + public String getHazelcastMapName() { + return hazelcastMapName; + } + + public boolean canProcess(String tableSDf) { + return this.tableSDf.equals(tableSDf); + } + + public Class getEntityClass() { + return entityClass; + } + + public static class TableProcessor { + public final Table inTable; + public final Function inGenerationId; + public final Function inSdfId; + + public TableProcessor(Table inTable, Function inGenerationId, Function inSdfId) { + this.inTable = inTable; + this.inGenerationId = inGenerationId; + this.inSdfId = inSdfId; + } + } +} diff --git a/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/service/PairSdfService.java b/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/service/PairSdfService.java new file mode 100644 index 000000000..46498016c --- /dev/null +++ b/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/service/PairSdfService.java @@ -0,0 +1,113 @@ +package ru.spcex.clearing.utility.service; + +import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.producer.Producer; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.InitializingBean; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.stereotype.Service; +import ru.clearing.classes.statics.data.register.PairSdf; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.platform.messaging.domain.BaseRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.utilities.PairSdfRequest; +import ru.spcex.clearing.platform.messaging.serialization.LogFormatter; +import ru.spcex.clearing.platform.messaging.service.QueueConsumer; +import ru.spcex.clearing.utility.Enum.Table; +import ru.spcex.platform.classes.base.SpcexObjectBase; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; + +import java.time.Instant; +import java.util.Arrays; +import java.util.Map; +import java.util.Optional; + +import static ru.spcex.clearing.platform.messaging.domain.Consts.PAIR_SDF; +import static ru.spcex.clearing.utility.Enum.Table.outTables; + +@Service +public class PairSdfService extends QueueConsumer implements InitializingBean { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final Imdg pairSdfImdg; + private final ImdgProvider imdgProvider; + + + @Autowired + public PairSdfService(Consumer kafkaQueue, Producer kafkaProducer, + ImdgProvider imdgProvider) { + super(kafkaQueue, kafkaProducer); + this.pairSdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_PairSdf, PairSdf.class); + this.imdgProvider = imdgProvider; + } + + @Override + public void afterPropertiesSet() { + callback(PairSdfRequest.class) + .setConsumer(this::processRequest) + .forDestination(PAIR_SDF, callbacks::put); + init(); + } + + public void processRequest(BaseRequest pairSdfNewRequestBaseRequest) { + log.info("Starting BaseRequest={} processing...", LogFormatter.toString(pairSdfNewRequestBaseRequest)); + PairSdfRequest request = pairSdfNewRequestBaseRequest.getRequestPayload(); + Optional table = Arrays.stream(Table.values()).filter(t -> t.canProcess(request.getTableSDf())).findFirst(); + if (table.isPresent()) { + if (outTables.containsKey(table.get())) { + processOut(request, table.get()); + } else { + processIn(request); + } + log.info("Request with BaseRequest.id={} successfully processed!", pairSdfNewRequestBaseRequest.getId()); + } else + log.info("Request with BaseRequest.id={} not processed, because not supported table = {}!", pairSdfNewRequestBaseRequest.getId(), request.getTableSDf()); + } + + + private void processIn(PairSdfRequest request) { + log.info("Process PairSdfRequest={} (by IN)!", LogFormatter.toString(request)); + Instant currDt = Instant.now(); + PairSdf pairSdf = new PairSdf(); + pairSdf.setCreated(currDt); + pairSdf.setUpdated(currDt); + pairSdf.setInSDfId(String.valueOf(request.getGenerationId())); + pairSdf.setInSDf(request.getFileNameSDf()); + pairSdfImdg.insert(pairSdf); + log.info("Successfully inserted PairSdf={} in map!", LogFormatter.toString(pairSdf)); + } + + private void processOut(PairSdfRequest request, Table outTable) { + log.info("Process PairSdfRequest={} (by OUT)!", LogFormatter.toString(request)); + Table.TableProcessor processor = outTables.get(outTable); + PairSdf pairSdf = null; + // out generationId -> outSdf inSdfId -> inSdf by id -> inSdf generationId -> pairSdf by in_s_df_id + if (request.getGenerationId() != null) { + Imdg outMap = imdgProvider.getImdg(outTable.getHazelcastMapName(), outTable.getEntityClass()); + SpcexObjectBase outSdf = outMap.getFirstObjectByFieldValues(Map.of("generationId", request.getGenerationId())); + if (outSdf != null && processor.inSdfId.apply(outSdf) != null) { + Imdg inMap = imdgProvider.getImdg(processor.inTable.getHazelcastMapName(), processor.inTable.getEntityClass()); + SpcexObjectBase inSdf = inMap.getSingleObjectByID(processor.inSdfId.apply(outSdf)); + if (inSdf != null && processor.inGenerationId.apply(inSdf) != null) { + pairSdf = pairSdfImdg.getFirstObjectByFieldValues(Map.of("inSDfId", processor.inGenerationId.apply(inSdf))); + } + } + } + Instant currDt = Instant.now(); + if (pairSdf != null) { + pairSdf.setUpdated(currDt); + pairSdf.setOutSDfId(String.valueOf(request.getGenerationId())); + pairSdf.setOutSDf(request.getFileNameSDf()); + pairSdfImdg.update(pairSdf); + log.info("Successfully updated PairSdf {} in map!", LogFormatter.toString(pairSdf)); + } else { + pairSdf = new PairSdf(); + pairSdf.setCreated(currDt); + pairSdf.setUpdated(currDt); + pairSdf.setOutSDfId(String.valueOf(request.getGenerationId())); + pairSdf.setOutSDf(request.getFileNameSDf()); + pairSdfImdg.insert(pairSdf); + log.info("Not found PairSdf(IN), inserted PairSdf={} in map!", LogFormatter.toString(pairSdf)); + } + } +} diff --git a/clearing-parent/utility-service/src/test/java/ru/spcex/clearing/utility/config/HazelcastInstanceTestConfiguration.java b/clearing-parent/utility-service/src/test/java/ru/spcex/clearing/utility/config/HazelcastInstanceTestConfiguration.java deleted file mode 100644 index 148e3e29b..000000000 --- a/clearing-parent/utility-service/src/test/java/ru/spcex/clearing/utility/config/HazelcastInstanceTestConfiguration.java +++ /dev/null @@ -1,23 +0,0 @@ -package ru.spcex.clearing.utility.config; - -import com.hazelcast.client.HazelcastClient; -import com.hazelcast.config.Config; -import com.hazelcast.config.JoinConfig; -import com.hazelcast.config.MulticastConfig; -import com.hazelcast.config.NetworkConfig; -import com.hazelcast.core.Hazelcast; -import com.hazelcast.core.HazelcastInstance; -import org.springframework.context.annotation.Bean; -import org.springframework.context.annotation.Configuration; - -import java.util.Random; - -@Configuration -public class HazelcastInstanceTestConfiguration { - - @Bean(name = "hazelcastInstance") - public HazelcastInstance hazelcastInstance() { - return HazelcastClient.newHazelcastClient(); - } - -} \ No newline at end of file diff --git a/clearing-parent/utility-service/src/test/java/ru/spcex/clearing/utility/config/HazelcastServiceTestConfiguration.java b/clearing-parent/utility-service/src/test/java/ru/spcex/clearing/utility/config/HazelcastServiceTestConfiguration.java deleted file mode 100644 index 6ec937bd0..000000000 --- a/clearing-parent/utility-service/src/test/java/ru/spcex/clearing/utility/config/HazelcastServiceTestConfiguration.java +++ /dev/null @@ -1,73 +0,0 @@ -package ru.spcex.clearing.utility.config; - -import com.hazelcast.config.*; -import com.hazelcast.core.Hazelcast; -import com.hazelcast.core.HazelcastInstance; -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.platform.imdg.iml.hazelcast.config.HazelcastClientParams; -import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService; -import ru.spcex.platform.imdg.iml.hazelcast.util.HazelcastHelper; - -import java.util.List; -import java.util.Random; - -@Configuration -public class HazelcastServiceTestConfiguration { - private HazelcastInstance hazelcastInstance; - - 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 = "hazelcastServiceTest") - public HazelcastService hazelcastService(@Qualifier("taskExecutorHazelcastClientInitializer") ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer, @Qualifier("taskExecutorIdGeneratorAwaiter") ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter, HazelcastClientParams params) { - Config cfg = new Config(); - cfg.setInstanceName("localhost"); - - NetworkConfig networkConfig = new NetworkConfig(); - JoinConfig joinConfig = new JoinConfig(); - joinConfig.setMulticastConfig(new MulticastConfig().setEnabled(false)); - joinConfig.setTcpIpConfig(new TcpIpConfig(). - setEnabled(true).setMembers(List.of("127.0.0.1"))); - networkConfig.setJoin(joinConfig); - - cfg.setNetworkConfig(networkConfig); - hazelcastInstance = Hazelcast.newHazelcastInstance(cfg); - - HazelcastHelper.imdgSystem_setStorageState(true, hazelcastInstance); - return new HazelcastService(taskExecutorHazelcastClientInitializer, - taskExecutorIdGeneratorAwaiter, params); - } - - @Bean(name = "taskExecutorHazelcastClientInitializer") - public ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer() { - return createThreadPoolTaskExecutor(1, true); - } - - @Bean(name = "taskExecutorIdGeneratorAwaiter") - public ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter() { - return createThreadPoolTaskExecutor(1, false); - } - - @Bean(name = "hazelcastClientParams") - public HazelcastClientParams getHazelcastClientParams() { - HazelcastClientParams params = new HazelcastClientParams(); - params.setLogin("dev"); - params.setPassword("dev-pass"); - params.setClusterMembers("127.0.0.1"); - params.setInstanceName("hzTestClient" + new Random().nextInt()); - params.setNearCacheConfig(new NearCacheConfig()); - - return params; - } -} \ No newline at end of file diff --git a/clearing-parent/utility-service/src/test/java/ru/spcex/clearing/utility/config/TestConfig.java b/clearing-parent/utility-service/src/test/java/ru/spcex/clearing/utility/config/TestConfig.java new file mode 100644 index 000000000..64f23c77b --- /dev/null +++ b/clearing-parent/utility-service/src/test/java/ru/spcex/clearing/utility/config/TestConfig.java @@ -0,0 +1,16 @@ +package ru.spcex.clearing.utility.config; + +import org.springframework.boot.context.properties.EnableConfigurationProperties; +import org.springframework.boot.test.context.TestConfiguration; +import org.springframework.context.annotation.ComponentScan; +import org.springframework.context.annotation.FilterType; +import ru.spcex.clearing.utility.config.settings.UtilityServiceSettings; + +@TestConfiguration +@ComponentScan(basePackages = {"ru.spcex.clearing.utility.config", + "ru.spcex.clearing.utility.service"}, + excludeFilters = {@ComponentScan.Filter(type = FilterType.ASSIGNABLE_TYPE, value = KafkaConfig.class), + @ComponentScan.Filter(type = FilterType.ASSIGNABLE_TYPE, value = UtilityServiceImdgConfig.class)}) +@EnableConfigurationProperties(value = UtilityServiceSettings.class) +public class TestConfig { +} \ No newline at end of file diff --git a/clearing-parent/utility-service/src/test/java/ru/spcex/clearing/utility/service/KeyRateServiceDisabled.java b/clearing-parent/utility-service/src/test/java/ru/spcex/clearing/utility/service/KeyRateServiceDisabled.java index dd7a3b93c..3fa7db0d5 100644 --- a/clearing-parent/utility-service/src/test/java/ru/spcex/clearing/utility/service/KeyRateServiceDisabled.java +++ b/clearing-parent/utility-service/src/test/java/ru/spcex/clearing/utility/service/KeyRateServiceDisabled.java @@ -15,8 +15,6 @@ import org.apache.kafka.common.TopicPartition; import org.junit.jupiter.api.Assertions; 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.misc.KeyRate; import ru.spcex.clearing.imdg.IMDGDistributedNames; @@ -26,8 +24,6 @@ import ru.spcex.clearing.platform.messaging.domain.Consts; import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteRequest; import ru.spcex.clearing.platform.messaging.domain.cud.utilities.KeyRateNewRequest; import ru.spcex.clearing.platform.messaging.domain.cud.utilities.KeyRateUpdateRequest; -import ru.spcex.clearing.utility.config.HazelcastInstanceTestConfiguration; -import ru.spcex.clearing.utility.config.HazelcastServiceTestConfiguration; import ru.spcex.clearing.utility.utils.MatcherFactory; import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService; @@ -37,9 +33,6 @@ import java.util.Collections; import java.util.HashMap; @ExtendWith(SpringExtension.class) -@ContextConfiguration(classes = { - HazelcastServiceTestConfiguration.class, - HazelcastInstanceTestConfiguration.class}) public class KeyRateServiceDisabled { static final ObjectMapper objectMapper = new ObjectMapper(); private final int TIMEOUT = 1000; @@ -47,10 +40,8 @@ public class KeyRateServiceDisabled { private MockConsumer kafkaMockQueue; private MockProducer mockProducer = new MockProducer<>(); @Autowired - @Qualifier("hazelcastServiceTest") private HazelcastService hazelcastService; @Autowired - @Qualifier("hazelcastInstance") private HazelcastInstance hz; private void setIsUsed(boolean isUsed) { diff --git a/clearing-parent/utility-service/src/test/java/ru/spcex/clearing/utility/service/PairSdfServiceTest.java b/clearing-parent/utility-service/src/test/java/ru/spcex/clearing/utility/service/PairSdfServiceTest.java new file mode 100644 index 000000000..fa8cb313c --- /dev/null +++ b/clearing-parent/utility-service/src/test/java/ru/spcex/clearing/utility/service/PairSdfServiceTest.java @@ -0,0 +1,103 @@ +package ru.spcex.clearing.utility.service; + +import org.apache.kafka.clients.consumer.MockConsumer; +import org.apache.kafka.clients.producer.Producer; +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.beans.factory.annotation.Qualifier; +import org.springframework.test.context.ContextConfiguration; +import org.springframework.test.context.junit.jupiter.SpringExtension; +import ru.clearing.classes.statics.data.register.PairSdf; +import ru.clearing.classes.statics.data.sdf.SDf01; +import ru.clearing.classes.statics.data.sdf.SDf02; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.platform.messaging.domain.cud.utilities.PairSdfRequest; +import ru.spcex.clearing.test.config.ImdgTestConfig; +import ru.spcex.clearing.test.config.KafkaTestConfig; +import ru.spcex.clearing.utility.config.TestConfig; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; + +import javax.annotation.PostConstruct; +import java.io.File; +import java.nio.file.Path; +import java.nio.file.Paths; + +import static ru.spcex.clearing.platform.messaging.domain.Consts.PAIR_SDF; +import static ru.spcex.clearing.test.TestUtils.*; +import static ru.spcex.clearing.test.config.ImdgTestConfig.currentID; +import static ru.spcex.clearing.test.config.ImdgTestConfig.waitAvailableImdgProviderAndAddAdminWithDefaultId; +import static ru.spcex.clearing.utility.Enum.Table.S_DF01; +import static ru.spcex.clearing.utility.Enum.Table.S_DF02; + +@ExtendWith(SpringExtension.class) +@ContextConfiguration(classes = {TestConfig.class, + ImdgTestConfig.class, + KafkaTestConfig.class}) +class PairSdfServiceTest { + private static final Long inGenerationId = currentID.getAndIncrement(); + private static final Long outGenerationId = currentID.getAndIncrement(); + private static final Long inSdfId = currentID.getAndIncrement(); + + @Autowired + private PairSdfService pairSdfService; + + @Autowired + @Qualifier("hazelcastServiceTest") + protected ImdgProvider imdgProvider; + + @Autowired + @Qualifier("mockProducer") + protected Producer mockProducer; + + private Imdg pairSdfImdg; + private Imdg sDf01Imdg; + private Imdg sDf02Imdg; + + static { + Path path = Paths.get("src", "main", "resources"); + String currentPath = path.toAbsolutePath().toString(); + System.setProperty("spring.config.location", currentPath + File.separator); + } + + @PostConstruct + public void init(){ + waitAvailableImdgProviderAndAddAdminWithDefaultId(); + this.sDf01Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf01, SDf01.class); + SDf01 sDf01 = new SDf01(); + sDf01.setId(inSdfId); + sDf01.setGenerationId(inGenerationId); + sDf01Imdg.insert(sDf01); + this.sDf02Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf02, SDf02.class); + SDf02 sDf02 = new SDf02(); + sDf02.setGenerationId(outGenerationId); + sDf02.setInSDfId(inSdfId); + sDf02Imdg.insert(sDf02); + this.pairSdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_PairSdf, PairSdf.class); + } + + @Disabled//Для отладки + @Test + void processRequest() throws InterruptedException { + Long ID = currentID.getAndIncrement(); + + PairSdfRequest pairSdfRequest = new PairSdfRequest(); + pairSdfRequest.setFileNameSDf("S_DF01_file_name"); + pairSdfRequest.setGenerationId(inGenerationId); + pairSdfRequest.setTableSDf(S_DF01.tableSDf); + + String jsonString = getJsonStringForNew(pairSdfRequest, ID); + addRecordToKafka((MockConsumer) pairSdfService.getConsumer(), PAIR_SDF, 0, 0, jsonString); + waitingSendAndCheckRecord(ID, mockProducer); + + pairSdfRequest.setFileNameSDf("S_DF02_file_name"); + pairSdfRequest.setGenerationId(outGenerationId); + pairSdfRequest.setTableSDf(S_DF02.tableSDf); + + jsonString = getJsonStringForNew(pairSdfRequest, ID); + addRecordToKafka((MockConsumer) pairSdfService.getConsumer(), PAIR_SDF, 0, 1, jsonString); + waitingSendAndCheckRecord(ID, mockProducer); + } +} \ No newline at end of file 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 a2fed7f78..b5327eb6d 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 @@ -115,6 +115,7 @@ public interface Consts { String ACCOUNT_NOTIFICATION_FEEDBACK = "account-notification-feedback"; String NOTIFICATION_NEW = "notification-new"; String NOTIFICATION_UPDATE = "notification-update"; + String PAIR_SDF = "pair-sdf"; String DESTINATION_SDF08_NEW = "s-df-08-new"; diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/utilities/PairSdfRequest.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/utilities/PairSdfRequest.java new file mode 100644 index 000000000..17935344a --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/utilities/PairSdfRequest.java @@ -0,0 +1,36 @@ +package ru.spcex.clearing.platform.messaging.domain.cud.utilities; + +import com.fasterxml.jackson.annotation.JsonProperty; + +public class PairSdfRequest { + @JsonProperty + public Long generationId; + @JsonProperty + public String fileNameSDf; + @JsonProperty + public String tableSDf; + + public Long getGenerationId() { + return generationId; + } + + public void setGenerationId(Long generationId) { + this.generationId = generationId; + } + + public String getFileNameSDf() { + return fileNameSDf; + } + + public void setFileNameSDf(String fileNameSDf) { + this.fileNameSDf = fileNameSDf; + } + + public String getTableSDf() { + return tableSDf; + } + + public void setTableSDf(String tableSDf) { + this.tableSDf = tableSDf; + } +}