From 29b3f3b1883010368495caa6b23b0647c886e632 Mon Sep 17 00:00:00 2001 From: ialbert Date: Tue, 30 May 2023 18:07:00 +0300 Subject: [PATCH] final MKR session duplicate intermediate MKR session --- .../clearing/service/EventsReceiver.java | 5 +- .../session/stage/FinalMkrSession.java | 228 ++++++++++++++++++ .../session/stage/IntermediateMkrSession.java | 2 +- .../session/stage/SessionManager.java | 5 +- .../platform/enumeration/SessionType.java | 2 +- 5 files changed, 238 insertions(+), 4 deletions(-) create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/FinalMkrSession.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 2aa318c28..a39ab9c5a 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 @@ -28,6 +28,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean { private final SecondaryAuctionT0Session secondaryAuctionT0Session; private final PrimaryAuctionB0Session primaryAuctionB0Session; private final IntermediateMkrSession intermediateMkrSession; + private final FinalMkrSession finalMkrSession; private final SessionManager sessionManager; @Autowired @@ -36,7 +37,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean { RegistryService registryService, PrimaryAuctionBnSession primaryAuctionBnSession, SecondaryAuctionT0Session secondaryAuctionT0Session, - PrimaryAuctionB0Session primaryAuctionB0Session, PrimaryAuctionT0Session primaryAuctionT0Session, IntermediateMkrSession intermediateMkrSession, SessionManager sessionManager) { + PrimaryAuctionB0Session primaryAuctionB0Session, PrimaryAuctionT0Session primaryAuctionT0Session, IntermediateMkrSession intermediateMkrSession, FinalMkrSession finalMkrSession, SessionManager sessionManager) { super(kafkaQueue); this.clearingService = clearingService; this.registryService = registryService; @@ -45,6 +46,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean { this.primaryAuctionB0Session = primaryAuctionB0Session; this.primaryAuctionT0Session = primaryAuctionT0Session; this.intermediateMkrSession = intermediateMkrSession; + this.finalMkrSession = finalMkrSession; this.sessionManager = sessionManager; } @@ -77,6 +79,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean { secondaryAuctionT0Session.continueSession(req); primaryAuctionB0Session.continueSession(req); intermediateMkrSession.continueSession(req); + finalMkrSession.continueSession(req); }) .forDestination(Consts.CONTINUE_SESSION_BN_FIRST_PART, callbacks::put); diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/FinalMkrSession.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/FinalMkrSession.java new file mode 100644 index 000000000..c62b2b938 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/FinalMkrSession.java @@ -0,0 +1,228 @@ +package ru.spcex.clearing.session.stage; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.InitializingBean; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.stereotype.Service; +import ru.clearing.classes.statics.data.execution.ExecutionCommon; +import ru.clearing.classes.statics.data.execution.ExecutionDeposit; +import ru.clearing.classes.statics.data.misc.Session; +import ru.clearing.classes.statics.data.payment.PaymentInstruction; +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.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.time.LocalDate; +import java.util.Collection; +import java.util.List; +import java.util.function.Function; +import java.util.function.Supplier; + +@Service +public class FinalMkrSession 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 InspectionObligations inspectionObligations; + private final FormingRegistersOnOS formingRegistersOnOS; + private final FormingPaymentInstruction formingPaymentInstruction; + private final UnlockResources unlockResources; + private final FinishingSession finishingSession; + private final EndStageNotification endStageNotification; + + private final Imdg executionDepositImdg; + private final Supplier> marketCodes; + + + public FinalMkrSession( + ImdgProvider imdgProvider, + BalanceRevise balanceRevise, + DealsPrepare dealsPrepare, + RequirementsAndObligationCreation requirementsAndObligationCreation, + ObligationAdmission obligationsAdmission, + InclusionObligations inclusionObligations, + FormingRegistersOnOS formingRegistersOnOS, + FormingPaymentInstruction formingPaymentInstruction, + UnlockResources unlockResources, + FinishingSession finishingSession, + EndStageNotification endStageNotification, + IMessageResolver messageResolver, + InspectionObligations inspectionObligations, + @Qualifier("marketCodesForBn") Supplier> marketCodes) { + 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.executionDepositImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionDeposit, ExecutionDeposit.class); + this.inspectionObligations = inspectionObligations; + this.marketCodes = marketCodes; + } + + @Override + public void afterPropertiesSet() throws Exception { + dealsPrepare.searchForExecutions(ExecutionType.ExecutionDeposit); + ImdgPredicateBuilder execFondPb = executionDepositImdg.predicateBuilder(); + dealsPrepare.addExecutionDepositCondition(execFondPb.regex("firstLegSettlementCode", "(^T0.*$)|(B[^0]\\d*$)")); +// dealsPrepare.addExecutionDepositCondition(execFondPb.in("market", marketCodes.get().toArray(new String[0]))); + } + + @Override + 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 (!isRunning()) { + return; + } + if (checkStage(TaskType.StartRevise)) { + firstPart(req); + } else if (checkStage(TaskType.FormingPaymentInstruction)) { + finishPart(req); + } else { + log.info("will not continue session, current stage is {}", currStage.get()); + } + } catch (StageException e) { + //already logged + } + } + + private void firstPart(BaseRequest req) { + try { + //stage 1 + StageResult> dealsPreparationResult; + { + DealsPreparePayload payload = new DealsPreparePayload(); + payload.setSessionId(currSession.getId()); + dealsPreparationResult = runStage(TaskType.DealsPrepare, payload, dealsPrepare); + //Сначала обрабатываются записи, у которых settlementCode = значению, у которого первый символ ="T", второй =0, + // а после settlementCode = значению, у которого первый символ ="B", второй ≠0, последующие =любые цифры + dealsPreparationResult.getStageResult().sort((o1, o2) -> { + String o1SettlementCode = ((ExecutionDeposit) o1).getFirstLegSettlementCode(); + String o2SettlementCode = ((ExecutionDeposit) o2).getSecondLegSettlementCode(); + Function mapper = (settlementCode) -> { + if (settlementCode.startsWith("T")) return -1; + else if (settlementCode.startsWith("B")) return 1; + else return 0; + }; + return mapper.apply(o1SettlementCode).compareTo(mapper.apply(o2SettlementCode)); + }); + } + //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, inspectionObligations); + } + //stage 6 + runStage(TaskType.FormingRegistersOnOS, formingRegistersOnOS); //returns Collection + //stage 7 + StageResult> paymentResult = runStage(TaskType.FormingPaymentInstruction, formingPaymentInstruction); + if (paymentResult.getStageResult().isEmpty()) { + runStage(TaskType.FormingPaymentInstruction, balanceRevise); +// finishPart(req); + } + } catch (StageException e) { + //already logged + } + } + + public void finishPart(BaseRequest req) { + try { + if (!checkStage(TaskType.FormingPaymentInstruction)) { + log.error("cannot continue session, current stage is {}", currStage.get()); + throw new StageException(); + } + //stage 10 + { + FinishingSessionPayload payload = new FinishingSessionPayload(); + payload.setSessionId(currSession.getId()); + runStage(TaskType.FinishingSession, payload, finishingSession); + } + //stage 11 + { + EndStageNotificationPayload payload = new EndStageNotificationPayload(); + payload.setSection(currSession.getSection()); + payload.setSessionId(currSession.getId()); + 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 { + TaskType startStatus = TaskType.StartRevise; + Session newSession = new Session(); + newSession.setSection(section().getKey()); + newSession.setSessionType(sectionType().getKey()); + newSession.setSessionStatus(startStatus.getKey()); + newSession.setWorkflowStatus(SessionStatus.ACTV.getKey()); + newSession.setClearingDate(LocalDate.now()); + + //todo companyId/securityId/userId передается из сообщения очереди + sessionImdg.insert(newSession); + currSession = newSession; + log.info("started new session.id={}", this.currSession.getId()); + currStage.set(startStatus); + return true; + } + } + } + + @Override + protected Section section() { + return Section.MKR; + } + + @Override + protected SessionType sectionType() { + return SessionType.FINL; + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/IntermediateMkrSession.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/IntermediateMkrSession.java index 7671c6e93..d829c632f 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/IntermediateMkrSession.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/IntermediateMkrSession.java @@ -84,7 +84,7 @@ public class IntermediateMkrSession extends AbstractSession implements Initializ dealsPrepare.searchForExecutions(ExecutionType.ExecutionDeposit); ImdgPredicateBuilder execFondPb = executionDepositImdg.predicateBuilder(); dealsPrepare.addExecutionDepositCondition(execFondPb.regex("firstLegSettlementCode", "(^T0.*$)|(B[^0]\\d*$)")); - dealsPrepare.addExecutionDepositCondition(execFondPb.in("market", marketCodes.get().toArray(new String[0]))); +// dealsPrepare.addExecutionDepositCondition(execFondPb.in("market", marketCodes.get().toArray(new String[0]))); } @Override diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/SessionManager.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/SessionManager.java index c75f721d8..9633ac4bb 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/SessionManager.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/SessionManager.java @@ -17,17 +17,19 @@ public class SessionManager { private final PrimaryAuctionB0Session primaryAuctionB0Session; private final SecondaryAuctionT0Session secondaryAuctionT0Session; private final IntermediateMkrSession intermediateMkrSession; + private final FinalMkrSession finalMkrSession; public SessionManager(PrimaryAuctionT0Session primaryAuctionT0Session, PrimaryAuctionBnSession primaryAuctionBnSession, PrimaryAuctionB0Session primaryAuctionB0Session, SecondaryAuctionT0Session secondaryAuctionT0Session, - IntermediateMkrSession intermediateMkrSession) { + IntermediateMkrSession intermediateMkrSession, FinalMkrSession finalMkrSession) { this.primaryAuctionT0Session = primaryAuctionT0Session; this.primaryAuctionBnSession = primaryAuctionBnSession; this.primaryAuctionB0Session = primaryAuctionB0Session; this.secondaryAuctionT0Session = secondaryAuctionT0Session; this.intermediateMkrSession = intermediateMkrSession; + this.finalMkrSession = finalMkrSession; } public void defineAndStartSession(BaseRequest r) { //task sclr section fond sessiontype trdt @@ -47,6 +49,7 @@ public class SessionManager { case IPO0 -> session = primaryAuctionB0Session; case TRDT -> session = secondaryAuctionT0Session; case MEDM -> session = intermediateMkrSession; + case FINL -> session = finalMkrSession; } if (session != null) { 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 c24aa9da3..249369a27 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"), IPO0("IPO0"), IPOT("IPOT"), TRDT("TRDT"), MEDM("MEDM"), + IPOB("IPOB"), IPO0("IPO0"), IPOT("IPOT"), TRDT("TRDT"), MEDM("MEDM"), FINL("FINL"), ; SessionType(String key) {