From c922351170132c17961c507800e2b60b828ba290 Mon Sep 17 00:00:00 2001 From: etreschenkov Date: Fri, 26 May 2023 21:17:46 +0300 Subject: [PATCH] fix statementService when we should send command for continue session --- .../config/MarketCodesBySessionConfig.java | 21 +++++++------------ .../clearing/service/EventsReceiver.java | 7 +++++-- .../clearing/service/StatementService.java | 8 +++---- .../service/executors/Sdf01Executor.java | 5 ----- .../session/stage/AbstractSession.java | 2 +- 5 files changed, 18 insertions(+), 25 deletions(-) diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/MarketCodesBySessionConfig.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/MarketCodesBySessionConfig.java index 5e8dd6bce..bffcb2729 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/MarketCodesBySessionConfig.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/MarketCodesBySessionConfig.java @@ -2,27 +2,22 @@ package ru.spcex.clearing.config; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; -import ru.clearing.classes.statics.data.misc.Market; -import ru.spcex.clearing.imdg.IMDGDistributedNames; -import ru.spcex.platform.enumeration.MarketType; -import ru.spcex.platform.enumeration.Section; -import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.ImdgProvider; -import java.util.Collection; import java.util.List; -import java.util.Map; -import java.util.stream.Collectors; @Configuration public class MarketCodesBySessionConfig { @Bean(name = "marketCodesForBn") public List marketCodesForBn(ImdgProvider imdgProvider) { - Imdg marketImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Market, Market.class); - Collection markets = marketImdg.getCollectionObjectsByFieldValues(Map.of( - "section", Section.FOND.getKey(), - "marketType", MarketType.PRMR.getKey())); - return markets.stream().map(Market::getCode).distinct().collect(Collectors.toList()); +// Imdg marketImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Market, Market.class); +// Collection markets = marketImdg.getCollectionObjectsByFieldValues(Map.of( +// "section", Section.FOND.getKey(), +// "marketType", MarketType.PRMR.getKey())); +// return markets.stream().map(Market::getCode).distinct().collect(Collectors.toList()); + return List.of("BMFC", "SMMC", "AEPC", "WBIC", "KBVC", "WBFC", "BMVC", "ABIC", + "AESC", "ABFC", "KBIC", "WBMC", "BMMC", "ABEC", "SMIC", "ABMC", "SMVC", "KBFC", "SKVC", "ABVC", "WEPC", + "KBMC", "SMFC", "BKVC", "WBVC", "BMIC", "WESC"); } } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/EventsReceiver.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/EventsReceiver.java index 05f9fe61e..e080ec3f7 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/EventsReceiver.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/EventsReceiver.java @@ -8,11 +8,14 @@ import ru.spcex.clearing.platform.messaging.domain.cud.clearing.CreateRegistryRe import ru.spcex.clearing.platform.messaging.domain.cud.clearing.Sdf04Request; import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonIdRequest; import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.utilities.STradesImportedRequest; import ru.spcex.clearing.platform.messaging.service.QueueConsumer; import ru.spcex.clearing.session.stage.PrimaryAuctionB0Session; import ru.spcex.clearing.session.stage.PrimaryAuctionBnSession; import ru.spcex.platform.enumeration.Task; +import static ru.spcex.clearing.platform.messaging.domain.Consts.S_TRADES_IMPORTED; + @Service public class EventsReceiver extends QueueConsumer implements InitializingBean { private final ClearingService clearingService; @@ -65,9 +68,9 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean { .setConsumer(primaryAuctionBnSession::continueSession) .forDestination(Consts.CONTINUE_SESSION_BN_FIRST_PART, callbacks::put); - callback(Object.class) + callback(STradesImportedRequest.class) .setConsumer(event -> clearingService.executeSTrade()) - .forDestination(Task.getOfTrades.topic(), callbacks::put); + .forDestination(S_TRADES_IMPORTED, callbacks::put); callback(CreateRegistryRequest.class) .setConsumer(event -> registryService.createRegistry(event.getRequestPayload())) .forDestination(Consts.REGISTRY_NEW, callbacks::put); diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java index e2b076037..17a544792 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java @@ -199,13 +199,13 @@ public class StatementService extends QueueConsumer implements InitializingBean Optional uncompletedPair = Optional.empty(); switch (sdfTable) { case SDF_01 -> - uncompletedPair = pairOfSdfRequest.entrySet().stream().filter(entry -> entry.getValue().getFirst() != null && entry.getValue().getSecond() == null).findFirst().map(Map.Entry::getKey); + uncompletedPair = pairOfSdfRequest.entrySet().stream().filter(entry -> entry.getValue().getFirst() == null && entry.getValue().getSecond() != null).findFirst().map(Map.Entry::getKey); case SDF_57 -> - uncompletedPair = pairOfSdfRequest.entrySet().stream().filter(entry -> entry.getValue().getFirst() == null && entry.getValue().getSecond() != null).findFirst().map(Map.Entry::getKey); - case SDF_04 -> uncompletedPair = pairOfSdfRequest.entrySet().stream().filter(entry -> entry.getValue().getFirst() != null && entry.getValue().getSecond() == null).findFirst().map(Map.Entry::getKey); - case SDF_13 -> + case SDF_04 -> uncompletedPair = pairOfSdfRequest.entrySet().stream().filter(entry -> entry.getValue().getFirst() == null && entry.getValue().getSecond() != null).findFirst().map(Map.Entry::getKey); + case SDF_13 -> + uncompletedPair = pairOfSdfRequest.entrySet().stream().filter(entry -> entry.getValue().getFirst() != null && entry.getValue().getSecond() == null).findFirst().map(Map.Entry::getKey); } return uncompletedPair; } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf01Executor.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf01Executor.java index 1d8f4bda9..cbb9e0fb7 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf01Executor.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf01Executor.java @@ -21,7 +21,6 @@ import ru.spcex.clearing.platform.messaging.domain.Consts; import ru.spcex.clearing.platform.messaging.domain.cud.account.sdf01.AccountSdfRequestPart; import ru.spcex.clearing.platform.messaging.domain.cud.balance.ExportToFileRequest; import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest; -import ru.spcex.clearing.platform.messaging.domain.cud.clearing.ContinueSessionBnRequest; import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.clearing.service.LoggingService; import ru.spcex.clearing.service.model.Result; @@ -95,10 +94,6 @@ public class Sdf01Executor extends AbstractExecutor { @Override public void sendCommand(KafkaSender kafkaSender, Result result) { - ContinueSessionBnRequest continueSessionBn = new ContinueSessionBnRequest(); - continueSessionBn.setGenerationId(result.getGenerationId()); - kafkaSender.sendRequestToQueue(Consts.CONTINUE_SESSION_BN_FIRST_PART, continueSessionBn); - ExportToFileRequest exportRequest = new ExportToFileRequest(); exportRequest.setSdfGroupId(result.getGenerationId()); exportRequest.setNameOfTable(exportTableName()); diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/AbstractSession.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/AbstractSession.java index c159463d4..a6c172cf7 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/AbstractSession.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/AbstractSession.java @@ -61,7 +61,7 @@ public class AbstractSession { protected boolean checkStage(TaskType t) { synchronized (this.currStage) { - return this.currStage.get().equals(t); + return t.equals(this.currStage.get()); } }