diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/SecondaryAuctionT0Session.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/SecondaryAuctionT0Session.java new file mode 100644 index 000000000..6df4b79b2 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/SecondaryAuctionT0Session.java @@ -0,0 +1,195 @@ +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.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.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.List; +import java.util.function.Supplier; + +@Service +public class SecondaryAuctionT0Session 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 executionFondImdg; + private final Supplier> marketCodes; + + + public SecondaryAuctionT0Session( + 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("marketCodesForT0") 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.executionFondImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionFond, ExecutionFond.class); + this.inspectionObligations = inspectionObligations; + this.marketCodes = marketCodes; + } + + @Override + public void afterPropertiesSet() throws Exception { + dealsPrepare.searchForExecutions(ExecutionType.ExecutionFond); + ImdgPredicateBuilder execFondPb = executionFondImdg.predicateBuilder(); + dealsPrepare.addExecutionFondCondition(execFondPb.regex("settlementCode", "^T0.*$")); + dealsPrepare.addExecutionFondCondition(execFondPb.in("market", marketCodes.get().toArray(new String[0]))); + } + + 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(); + } + if (checkStage(TaskType.StartRevise)) { + firstPart(req); + } else if (checkStage(TaskType.FormingPaymentInstruction)) { + finishPart(req); + } + } 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); + } + //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 + runStage(TaskType.FormingPaymentInstruction, formingPaymentInstruction); + } 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()); + 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.TRDT.getKey()); + newSession.setSessionStatus(SessionStatus.CLRN.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(TaskType.StartRevise); + return true; + } + } + } +} 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 b299f910f..11fb75871 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"), + IPOB("IPOB"), IPO0("IPO0"), IPOT("IPOT"), TRDT("TRDT"), ; SessionType(String key) {