From 936e4613d8af210bb0cc0a629802a028faa5b010 Mon Sep 17 00:00:00 2001 From: ialbert Date: Tue, 23 May 2023 14:43:43 +0300 Subject: [PATCH] http://jira.mfd.msk:8088/browse/CLS-290 --- .../stage/PrimaryAuctionBnSession.java | 180 ++++++++++++++---- .../session/stage/impl/DealsPrepare.java | 10 +- 2 files changed, 152 insertions(+), 38 deletions(-) diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/PrimaryAuctionBnSession.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/PrimaryAuctionBnSession.java index d98c37db4..507d1014a 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/PrimaryAuctionBnSession.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/PrimaryAuctionBnSession.java @@ -3,12 +3,17 @@ package ru.spcex.clearing.session.stage; import org.apache.kafka.clients.consumer.Consumer; 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.misc.Session; import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.platform.messaging.domain.BaseRequest; +import ru.spcex.clearing.platform.messaging.domain.Consts; import ru.spcex.clearing.platform.messaging.service.QueueConsumer; import ru.spcex.clearing.session.stage.impl.*; +import ru.spcex.clearing.session.stage.task.*; import ru.spcex.platform.enumeration.Section; import ru.spcex.platform.enumeration.SessionStatus; import ru.spcex.platform.enumeration.SessionType; @@ -16,10 +21,11 @@ import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.ImdgProvider; import ru.spcex.platform.utils.enumeration.IMessageResolver; -import java.util.concurrent.atomic.AtomicBoolean; +import java.util.List; +import java.util.concurrent.atomic.AtomicReference; @Service -public class PrimaryAuctionBnSession extends QueueConsumer { +public class PrimaryAuctionBnSession extends QueueConsumer implements InitializingBean { private final Logger log = LoggerFactory.getLogger(getClass()); private final BalanceRevise balanceRevise; private final DealsPrepare dealsPrepare; @@ -30,13 +36,13 @@ public class PrimaryAuctionBnSession extends QueueConsumer { private final FormingRegistersOnOS formingRegistersOnOS; private final FormingPaymentInstruction formingPaymentInstruction; private final UnlockResources unlockResources; + private final FinishingSession finishingSession; private final EndStageNotification endStageNotification; - // - private final AtomicBoolean running; - private Long sessionId; private Imdg sessionImdg; private final IMessageResolver messageResolver; + private final AtomicReference currStage = new AtomicReference<>(); + private Session currSession; public PrimaryAuctionBnSession( @Qualifier("createConsumer") Consumer kafkaQueue, @@ -49,7 +55,7 @@ public class PrimaryAuctionBnSession extends QueueConsumer { FormingRegistersOnOS formingRegistersOnOS, FormingPaymentInstruction formingPaymentInstruction, UnlockResources unlockResources, - EndStageNotification endStageNotification, IMessageResolver messageResolver) { + FinishingSession finishingSession, EndStageNotification endStageNotification, IMessageResolver messageResolver) { super(kafkaQueue); this.balanceRevise = balanceRevise; this.dealsPrepare = dealsPrepare; @@ -59,43 +65,151 @@ public class PrimaryAuctionBnSession extends QueueConsumer { this.formingRegistersOnOS = formingRegistersOnOS; this.formingPaymentInstruction = formingPaymentInstruction; this.unlockResources = unlockResources; + this.finishingSession = finishingSession; this.endStageNotification = endStageNotification; - // + this.sessionImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Session, Session.class); this.messageResolver = messageResolver; - this.running = new AtomicBoolean(false); + } + + @Override + public void afterPropertiesSet() throws Exception { + callback(Object.class) + .setConsumer(this::continueSession) + .forDestination(Consts.SDF57_PROCESS, callbacks::put); + init(); } public void runSession() { - synchronized (this.running) { - if (this.running.get()) { - log.info("already running session.id={}", this.sessionId); - return; + 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, companyIdPayload, 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 StageResult runStage(TaskType type, ISessionStage stage) { + return runStage(type, null, stage); + } + + @SuppressWarnings("unchecked") + private StageResult runStage(TaskType type, T payload, ISessionStage stage) { + continueRunning(type); + log.info("session.id={} step {} started", currSession.getId(), currStage.get()); + Task t = new Task<>(type, payload); + StageResult stgRes = stage.submit(t); + log.info("session.id={} step {} result: {} ", + currSession.getId(), + currStage.get(), + stgRes.success ? "success" : messageResolver.resolve(stgRes.error)); + if (!stgRes.success) { + endSession(); + throw new StageException(); + } + return (StageResult) stgRes; + } + + private boolean startSession() { + synchronized (this.currStage) { + if (this.currStage.get() != null) { + log.info("already running session.id={}", this.currSession.getId()); + return false; } else { - this.running.set(true); + Session newSession = new Session(); + newSession.setSection(Section.FOND.getKey()); + newSession.setSessionType(SessionType.IPOB.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; } } - Session newSession = new Session(); - newSession.setSection(Section.FOND.getKey()); - newSession.setSessionType(SessionType.IPOB.getKey()); - newSession.setSessionStatus(SessionStatus.CLRN.getKey()); - //todo companyId/securityId/userId передается из сообщения очереди - sessionImdg.insert(newSession); - sessionId = newSession.getId(); - //java.util.function.Consumer> logError = stageResult -> { - // if (!stageResult.success) { - // log.error("stageResult.error={}", messageResolver.resolve(stageResult.error)); - // newSession.setSessionStatus(SessionStatus.CLOS.getKey()); - // } - //}; - //{ - // StageResult reviseResult = balanceRevise.submit(task(TaskType.StartRevise)); - // logError.accept(reviseResult); - //} - } - private Task task(TaskType taskType) { - return new Task<>(taskType, null); + private void endSession() { + synchronized (this.currStage) { + this.currStage.set(null); + this.currSession = null; + } + } + + private void continueRunning(TaskType t) { + synchronized (this.currStage) { + this.currStage.set(t); + } + } + + private boolean checkStage(TaskType t) { + synchronized (this.currStage) { + return this.currStage.get().equals(t); + } + } + + private static class StageException extends RuntimeException { } } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/DealsPrepare.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/DealsPrepare.java index ab25bb192..d0613654c 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/DealsPrepare.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/DealsPrepare.java @@ -4,6 +4,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Autowired; 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.execution.ExecutionFond; import ru.spcex.clearing.imdg.IMDGDistributedNames; @@ -11,7 +12,6 @@ import ru.spcex.clearing.session.stage.ISessionStage; import ru.spcex.clearing.session.stage.StageResult; import ru.spcex.clearing.session.stage.Task; import ru.spcex.clearing.session.stage.task.DealsPreparePayload; -import ru.spcex.platform.classes.base.interfaces.IExecution; import ru.spcex.platform.classes.base.interfaces.WithExchangeExecutionId; import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.ImdgProvider; @@ -74,11 +74,11 @@ public class DealsPrepare implements ISessionStage { ImdgPredicate excFondPrct = prdComposer.apply(execFondPredicates, executionFondImdg); Collection excDpsts = executionDepositImdg.getCollectionObjectsByPredicate(excDepPrct); Collection excFonds = executionFondImdg.getCollectionObjectsByPredicate(excFondPrct); - List excs = Stream.concat(excDpsts.stream().map(execToInterface()), + List excs = Stream.concat(excDpsts.stream().map(execToInterface()), excFonds.stream().map(execToInterface())) .sorted(Comparator.comparing(WithExchangeExecutionId::getExchangeExecutionId)) .toList(); - for (IExecution exc : excs) { + for (ExecutionCommon exc : excs) { exc.setSessionId(sessionId); if (exc instanceof ExecutionDeposit) { executionDepositImdg.update((ExecutionDeposit) exc); @@ -86,7 +86,7 @@ public class DealsPrepare implements ISessionStage { executionFondImdg.update((ExecutionFond) exc); } } - StageResult> res = new StageResult<>(null, true); + StageResult> res = new StageResult<>(null, true); res.setStageResult(excs); return res; } @@ -103,7 +103,7 @@ public class DealsPrepare implements ISessionStage { - private static Function execToInterface() { + private static Function execToInterface() { return (e) -> e; } }