From 0fe495a1a2649dea1720c794fe170866cb9dd425 Mon Sep 17 00:00:00 2001 From: ialbert Date: Wed, 24 May 2023 15:19:10 +0300 Subject: [PATCH] B0 session --- .../clearing/service/EventsReceiver.java | 22 ++- .../stage/PrimaryAuctionB0Session.java | 165 ++++++++++++++++++ .../platform/enumeration/MarketType.java | 18 ++ .../platform/enumeration/SessionType.java | 2 +- .../ru/spcex/platform/enumeration/Task.java | 1 + 5 files changed, 203 insertions(+), 5 deletions(-) create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/PrimaryAuctionB0Session.java create mode 100644 platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/MarketType.java 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 758b9be7b..3e51ed48e 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 @@ -9,6 +9,7 @@ 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.service.QueueConsumer; +import ru.spcex.clearing.session.stage.PrimaryAuctionB0Session; import ru.spcex.clearing.session.stage.PrimaryAuctionBnSession; import ru.spcex.platform.enumeration.Task; @@ -17,15 +18,18 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean { private final ClearingService clearingService; private final RegistryService registryService; private final PrimaryAuctionBnSession primaryAuctionBnSession; + private final PrimaryAuctionB0Session primaryAuctionB0Session; public EventsReceiver(Consumer kafkaQueue, ClearingService clearingService, RegistryService registryService, - PrimaryAuctionBnSession primaryAuctionBnSession) { + PrimaryAuctionBnSession primaryAuctionBnSession, + PrimaryAuctionB0Session primaryAuctionB0Session) { super(kafkaQueue); this.clearingService = clearingService; this.registryService = registryService; this.primaryAuctionBnSession = primaryAuctionBnSession; + this.primaryAuctionB0Session = primaryAuctionB0Session; } @Override @@ -42,15 +46,25 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean { callback(LauncherCommandRequest.class) .setConsumer(event -> clearingService.executeVerification()) .forDestination(Task.getVerification.topic(), callbacks::put); - callback(Object.class) - .setConsumer(primaryAuctionBnSession::runSession) - .forDestination(Task.startOfClearing.topic(), callbacks::put); callback(CommonIdRequest.class) .setConsumer(clearingService::continueClearing) .forDestination(Consts.CONTINUE_CLEARING, callbacks::put); + + callback(Object.class) + .setConsumer(primaryAuctionB0Session::runSession) + .forDestination(Task.startOfB0.topic(), callbacks::put); +// callback(Object.class) +// .setConsumer(primaryAuctionB0Session::continueSession) +// .forDestination(Task.startOfB0.topic(), callbacks::put); + + + callback(Object.class) + .setConsumer(primaryAuctionBnSession::runSession) + .forDestination(Task.startOfClearing.topic(), callbacks::put); callback(Object.class) .setConsumer(primaryAuctionBnSession::continueSession) .forDestination(Consts.SDF57_PROCESS, callbacks::put); + callback(Object.class) .setConsumer(event -> clearingService.executeSTrade()) .forDestination(Task.getOfTrades.topic(), callbacks::put); diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/PrimaryAuctionB0Session.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/PrimaryAuctionB0Session.java new file mode 100644 index 000000000..543bb940e --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/PrimaryAuctionB0Session.java @@ -0,0 +1,165 @@ +package ru.spcex.clearing.session.stage; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.InitializingBean; +import org.springframework.stereotype.Service; +import ru.clearing.classes.statics.data.execution.ExecutionCommon; +import ru.clearing.classes.statics.data.execution.ExecutionFond; +import ru.clearing.classes.statics.data.misc.Session; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.platform.messaging.domain.BaseRequest; +import ru.spcex.clearing.session.stage.impl.*; +import ru.spcex.clearing.session.stage.task.*; +import ru.spcex.platform.classes.base.interfaces.ExecutionType; +import ru.spcex.platform.enumeration.MarketType; +import ru.spcex.platform.enumeration.Section; +import ru.spcex.platform.enumeration.SessionStatus; +import ru.spcex.platform.enumeration.SessionType; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; +import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder; +import ru.spcex.platform.utils.enumeration.IMessageResolver; + +import java.util.List; + +@Service +public class PrimaryAuctionB0Session extends AbstractSession implements InitializingBean { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final BalanceRevise balanceRevise; + private final DealsPrepare dealsPrepare; + private final RequirementsAndObligationCreation requirementsAndObligationCreation; + private final ObligationAdmission obligationsAdmission; + + private final InclusionObligations inclusionObligations; + private final FormingRegistersOnOS formingRegistersOnOS; + private final FormingPaymentInstruction formingPaymentInstruction; + private final UnlockResources unlockResources; + private final FinishingSession finishingSession; + private final EndStageNotification endStageNotification; + + private final Imdg executionFondImdg; + + public PrimaryAuctionB0Session( + ImdgProvider imdgProvider, + BalanceRevise balanceRevise, + DealsPrepare dealsPrepare, + RequirementsAndObligationCreation requirementsAndObligationCreation, + ObligationAdmission obligationsAdmission, + InclusionObligations inclusionObligations, + FormingRegistersOnOS formingRegistersOnOS, + FormingPaymentInstruction formingPaymentInstruction, + UnlockResources unlockResources, + FinishingSession finishingSession, EndStageNotification endStageNotification, IMessageResolver messageResolver) { + super(imdgProvider, messageResolver); + this.balanceRevise = balanceRevise; + this.dealsPrepare = dealsPrepare; + this.requirementsAndObligationCreation = requirementsAndObligationCreation; + this.obligationsAdmission = obligationsAdmission; + this.inclusionObligations = inclusionObligations; + this.formingRegistersOnOS = formingRegistersOnOS; + this.formingPaymentInstruction = formingPaymentInstruction; + this.unlockResources = unlockResources; + this.finishingSession = finishingSession; + this.endStageNotification = endStageNotification; + this.executionFondImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionFond, ExecutionFond.class); + } + + @Override + public void afterPropertiesSet() throws Exception { + dealsPrepare.searchForExecutions(ExecutionType.ExecutionFond); + ImdgPredicateBuilder execFondPb = executionFondImdg.predicateBuilder(); + dealsPrepare.addExecutionFondCondition(execFondPb.regex("settlementCode", "^B0.*$")); + dealsPrepare.addExecutionFondCondition(execFondPb.equals("marketType", MarketType.PRMR.getKey())); + } + + public void runSession(BaseRequest req) { + if (!startSession()) { + return; + } + StageResult submit = balanceRevise.submit(new Task<>(TaskType.StartRevise, null)); + if (!submit.success) { + log.error("stage {} error={}", balanceRevise.getClass().getSimpleName(), messageResolver.resolve(submit.error)); + endSession(); + } else { + log.info("stage BalanceRevise success, waiting for a response from kafka"); + } + } + + public void continueSession(BaseRequest req) { + try { + if (!checkStage(TaskType.StartRevise)) { + log.error("cannot continue session, current stage is {}", currStage.get()); + throw new StageException(); + } + //stage 0 + runStage(TaskType.ContinueRevise, balanceRevise); + //stage 1 + StageResult> dealsPreparationResult; + { + DealsPreparePayload payload = new DealsPreparePayload(); + payload.setSessionId(currSession.getId()); + dealsPreparationResult = runStage(TaskType.DealsPrepare , payload, dealsPrepare); + } + //stage 2 + runStage(TaskType.RequirementsAndObligationsCreate, dealsPreparationResult.getStageResult(), requirementsAndObligationCreation); + //stage 3 + runStage(TaskType.ObligationsAdmission, currSession.getId(), obligationsAdmission); + //stage 4 + { + InclusionToPoolPayload inclusionToPoolPayload = new InclusionToPoolPayload(); + inclusionToPoolPayload.setSessionType(currSession.getSessionType()); + runStage(TaskType.InclusionToPool, inclusionToPoolPayload, inclusionObligations); + } + //stage 5 + { + InspectionPoolPayload companyIdPayload = new InspectionPoolPayload(); + companyIdPayload.setProcessedCompanyId(currSession.getCompanyId()); + runStage(TaskType.InspectionObligations, companyIdPayload, inclusionObligations); + } + //stage 6 + runStage(TaskType.FormingRegistersOnOS, formingRegistersOnOS); //returns Collection + //stage 7 + runStage(TaskType.FormingPaymentInstruction, formingPaymentInstruction); + //stage 8 + { + UnlockResourcesPayload unlockResourcesPayload = new UnlockResourcesPayload(); + //todo set arguments + runStage(TaskType.UnlockResources, unlockResourcesPayload, unlockResources); //returns Collection + } + //stage 9 + { + FinishingSessionPayload payload = new FinishingSessionPayload(); + payload.setSessionId(currSession.getId()); + runStage(TaskType.FinishingSession, payload, finishingSession); + } + { + EndStageNotificationPayload payload = new EndStageNotificationPayload(); + payload.setSection(currSession.getSection()); + runStage(TaskType.EndStageNotification, payload, endStageNotification); + } + } catch (StageException e) { + //already logged + } + } + + private boolean startSession() { + synchronized (this.currStage) { + if (this.currStage.get() != null) { + log.info("already running session.id={}", this.currSession.getId()); + return false; + } else { + Session newSession = new Session(); + newSession.setSection(Section.FOND.getKey()); + newSession.setSessionType(SessionType.IPO0.getKey()); + newSession.setSessionStatus(SessionStatus.CLRN.getKey()); + //todo companyId/securityId/userId передается из сообщения очереди + sessionImdg.insert(newSession); + currSession = newSession; + log.info("started new session.id={}", this.currSession.getId()); + currStage.set(TaskType.StartRevise); + return true; + } + } + } +} diff --git a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/MarketType.java b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/MarketType.java new file mode 100644 index 000000000..930e40a2c --- /dev/null +++ b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/MarketType.java @@ -0,0 +1,18 @@ +package ru.spcex.platform.enumeration; + +import ru.spcex.platform.utils.enumeration.IEnumKey; + +public enum MarketType implements IEnumKey { + PRMR("PRMR"); + + private final String key; + + MarketType(String key) { + this.key = key; + } + + @Override + public String getKey() { + return key; + } +} diff --git a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/SessionType.java b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/SessionType.java index 70ba05bd5..43ea888fb 100644 --- a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/SessionType.java +++ b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/SessionType.java @@ -3,7 +3,7 @@ package ru.spcex.platform.enumeration; import ru.spcex.platform.utils.enumeration.IEnumKey; public enum SessionType implements IEnumKey { - IPOB("IPOB"), + IPOB("IPOB"), IPO0("IPO0"), ; SessionType(String key) { diff --git a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/Task.java b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/Task.java index 40bd8c007..422a649ae 100644 --- a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/Task.java +++ b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/Task.java @@ -10,6 +10,7 @@ public enum Task implements IEnumKey { getVerification("GVER"),// Запуск сверки @Deprecated /* todo GBLD удаляется по CLS-267, CLS-275 */ getBalance("GBLD"),// Поступление средств startOfClearing("SCLR"),// Запуск клиринговой сессии + startOfB0("IPO0"),// Запуск клиринговой сессии startOfPreClearing("SPRC"),// Запуск преклиринга startPostClearing("SPOC"),// Запуск постклиринга createOrder("CORD"),