From 33bc8ce75df5f4486612878e8fdfb3ef472c7a9f Mon Sep 17 00:00:00 2001 From: ialbert Date: Tue, 1 Oct 2024 13:14:32 +0300 Subject: [PATCH] PAYM session --- .../clearing/session/stage/PaymSession.java | 251 ++++++++++++++++++ .../impl/FormingPaymentInstructionAssets.java | 4 +- .../stage/impl/InclusionObligations.java | 7 +- .../stage/impl/InspectionObligationsV2.java | 2 +- 4 files changed, 260 insertions(+), 4 deletions(-) create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/PaymSession.java diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/PaymSession.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/PaymSession.java new file mode 100644 index 000000000..e12e8b58e --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/PaymSession.java @@ -0,0 +1,251 @@ +package ru.spcex.clearing.session.stage; + +import java.time.Instant; +import java.time.LocalDate; +import java.util.List; +import java.util.function.Supplier; +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.ExecutionCurrency; +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.BalanceRevise; +import ru.spcex.clearing.session.stage.impl.DealsPrepare; +import ru.spcex.clearing.session.stage.impl.EndStageNotification; +import ru.spcex.clearing.session.stage.impl.FinishingSession; +import ru.spcex.clearing.session.stage.impl.FormingPaymentInstructionAssets; +import ru.spcex.clearing.session.stage.impl.FormingRegistersOnOS; +import ru.spcex.clearing.session.stage.impl.InclusionObligations; +import ru.spcex.clearing.session.stage.impl.InspectionObligationsV2; +import ru.spcex.clearing.session.stage.impl.ObligationAdmission; +import ru.spcex.clearing.session.stage.impl.PaymentInfo; +import ru.spcex.clearing.session.stage.impl.RequirementsAndObligationCreation; +import ru.spcex.clearing.session.stage.impl.UnlockResources; +import ru.spcex.clearing.session.stage.monitor.SessionMonitor; +import ru.spcex.clearing.session.stage.monitor.SessionMonitorFactory; +import ru.spcex.clearing.session.stage.task.DealsPreparePayload; +import ru.spcex.clearing.session.stage.task.EndStageNotificationPayload; +import ru.spcex.clearing.session.stage.task.FinishingSessionPayload; +import ru.spcex.clearing.session.stage.task.FormingPaymentInstructionPayload; +import ru.spcex.clearing.session.stage.task.InclusionToPoolPayload; +import ru.spcex.clearing.session.stage.task.InspectionPoolPayload; +import ru.spcex.clearing.session.stage.task.RequirementsAndObligationCreationPayload; +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; + +@Service +public class PaymSession extends AbstractSession implements InitializingBean { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final DealsPrepare dealsPrepare; + private final RequirementsAndObligationCreation requirementsAndObligationCreation; + private final ObligationAdmission obligationsAdmission; + private final BalanceRevise balanceRevise; + private final InclusionObligations inclusionObligations; + + private final InspectionObligationsV2 inspectionObligations; + private final FormingPaymentInstructionAssets formingPaymentInstructionAssets; + private final UnlockResources unlockResources; + private final FinishingSession finishingSession; + private final EndStageNotification endStageNotification; + + private final Imdg executionFondImdg; + private final Imdg executionCurrencyImdg; + private final Supplier> marketCodes; + private SessionMonitor afterPaymentsReviseMonitor; + private SessionMonitor afterReviseErrorMonitor; + + + public PaymSession( + ImdgProvider imdgProvider, + DealsPrepare dealsPrepare, + RequirementsAndObligationCreation requirementsAndObligationCreation, + ObligationAdmission obligationsAdmission, + InclusionObligations inclusionObligations, + FormingRegistersOnOS formingRegistersOnOS, BalanceRevise balanceRevise, + FormingPaymentInstructionAssets formingPaymentInstructionAssets, + UnlockResources unlockResources, + FinishingSession finishingSession, + EndStageNotification endStageNotification, + IMessageResolver messageResolver, + InspectionObligationsV2 inspectionObligations, + @Qualifier("marketCodesForCurr") Supplier> marketCodes) { + super(imdgProvider, messageResolver); + this.dealsPrepare = dealsPrepare; + this.requirementsAndObligationCreation = requirementsAndObligationCreation; + this.obligationsAdmission = obligationsAdmission; + this.inclusionObligations = inclusionObligations; + this.balanceRevise = balanceRevise; + this.formingPaymentInstructionAssets = formingPaymentInstructionAssets; + this.unlockResources = unlockResources; + this.finishingSession = finishingSession; + this.endStageNotification = endStageNotification; + this.executionFondImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionFond, ExecutionFond.class); + this.executionCurrencyImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionCurrency, ExecutionCurrency.class); + this.inspectionObligations = inspectionObligations; + this.marketCodes = marketCodes; + } + + @Override + public void afterPropertiesSet() throws Exception { + dealsPrepare.searchForExecutions(ExecutionType.ExecutionCurrency); + ImdgPredicateBuilder execFondPb = executionCurrencyImdg.predicateBuilder(); + dealsPrepare.addExecutionCurrencyCondition(execFondPb.in("market", marketCodes.get().toArray(new String[0]))); + formingPaymentInstructionAssets.setSessionType(sessionType()); + inclusionObligations.setSessionType(sessionType()); + + inspectionObligations.setSection(section()); + finishingSession.setSection(section()); + finishingSession.setSessionType(sessionType()); + imdgProvider.waitAvailable(); + initSessionIfPresent(); + } + + @Override + public void runSession(BaseRequest req) { + if (!startSession()) { + return; + } + try { + //stage 1 + StageResult> dealsPreparationResult; + { + DealsPreparePayload payload = new DealsPreparePayload(); + payload.setSessionId(currSession.getId()); + dealsPreparationResult = runStage(TaskType.DealsPrepare, payload, dealsPrepare); + } + { + //stage 2 + RequirementsAndObligationCreationPayload payload = new RequirementsAndObligationCreationPayload( + dealsPreparationResult.getStageResult(), currSession.getId() + ); + runStage(TaskType.RequirementsAndObligationsCreate, payload, requirementsAndObligationCreation); + } + //stage 3 + runStage(TaskType.ObligationsAdmission, currSession.getId(), obligationsAdmission); + //stage 4 + { + InclusionToPoolPayload inclusionToPoolPayload = new InclusionToPoolPayload(); + inclusionToPoolPayload.setSessionType(currSession.getSessionType()); + inclusionToPoolPayload.setSessionId(currSession.getId()); + runStage(TaskType.InclusionToPool, inclusionToPoolPayload, inclusionObligations); + } + //stage 5 + { + InspectionPoolPayload companyIdPayload = new InspectionPoolPayload(); + companyIdPayload.setSessionId(currSession.getId()); + companyIdPayload.setProcessedCompanyId(currSession.getCompanyId()); + runStage(TaskType.InspectionObligations, companyIdPayload, inspectionObligations); + } + //stage 7 + StageResult paymentResult = null; + { + FormingPaymentInstructionPayload payload = new FormingPaymentInstructionPayload(); + payload.setSessionId(currSession.getId()); + paymentResult = runStage(TaskType.FormingPaymentInstruction, payload, formingPaymentInstructionAssets); + } + if (paymentResult.getStageResult().getPaymentInstructions().isEmpty()) { + log.info("no payment instructions were created"); + finishPart(); + } else { + runStage(TaskType.FormingPaymentInstruction, balanceRevise); + this.afterPaymentsReviseMonitor = SessionMonitorFactory.waitRevise(); + log.info("created {} PaymentInstructions, waiting for {}", + paymentResult.getStageResult().getPaymentInstructions().size(), + this.afterPaymentsReviseMonitor.allConditions()); + } + } catch (StageException e) { + //already logged + } + } + + public void continueSession(BaseRequest req) { + try { + if (!isRunning()) { + return; + } + log.info("session is running, stage {}, monitors: {}", + currStage.get(), + logMonitors(afterPaymentsReviseMonitor)); + if (afterReviseErrorMonitor != null && isMonitorPassed(afterReviseErrorMonitor, req.getRequestPayload())) { + afterReviseErrorMonitor = null; + finishPart(); + return; + } + } catch (StageException e) { + //already logged + } + } + + public void finishPart() { + try { + //stage 9 continue revision + { + runStage(TaskType.ContinueRevise, currSession.getId(), balanceRevise, false); + } + //stage 10 + { + FinishingSessionPayload payload = new FinishingSessionPayload(); + payload.setSessionId(currSession.getId()); + payload.setPr("1"); + runStage(TaskType.FinishingSession, payload, finishingSession); + } + //stage 11 + { + EndStageNotificationPayload payload = new EndStageNotificationPayload(); + payload.setSection(currSession.getSection()); + payload.setSessionId(currSession.getId()); + runStage(TaskType.EndStageNotification, payload, endStageNotification); + endSession(); + } + } 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.MKR.getKey()); + newSession.setSessionType(sessionType().getKey()); + newSession.setSessionStatus(startStatus.getKey()); + newSession.setWorkflowStatus(SessionStatus.ACTV.getKey()); + newSession.setClearingDate(LocalDate.now()); + newSession.setCreated(Instant.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; + } + } + } + + @Override + protected Section section() { + return Section.MKR; + } + + @Override + protected SessionType sessionType() { + return SessionType.PAYM; + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/FormingPaymentInstructionAssets.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/FormingPaymentInstructionAssets.java index 6cfbcfab2..532e9076b 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/FormingPaymentInstructionAssets.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/FormingPaymentInstructionAssets.java @@ -407,7 +407,7 @@ public class FormingPaymentInstructionAssets implements ISessionStage { Registry rgs = pair.getSecond(); Account account = accountImdg.getSingleObjectByID(paymentInstruction.getCreditLeg_accountId()); Security security = securityImdg.getSingleObjectByID(paymentInstruction.getCreditLeg_securityId()); - if (SessionType.CURR.equals(sessionType) || (rgs != null + if (SessionType.CURR.equals(sessionType) || SessionType.PAYM.equals(sessionType) || (rgs != null && RegistryInstrumentType.M.equalsByKey(rgs.getRegistryInstrumentType()) && !CurrencyCode.isRub(rgs.getSecuritySymbol()))) { boolean madeStatements = statementsWereMade(paymentInstruction, rgs); @@ -496,7 +496,7 @@ public class FormingPaymentInstructionAssets implements ISessionStage { private boolean statementsWereMade(PaymentInstruction pmt, Registry am_b) { ImdgPredicateBuilder pb = stlmHPropsImdg.predicateBuilder(); String curCode = pmt.getCreditLeg_currencyCode(); - { + if (!SessionType.PAYM.equals(sessionType)) { SettlementHouseProperties sttHs = stlmHPropsImdg.getFirstObjectByPredicate( pb.and( pb.in("currencyCode", curCode), diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InclusionObligations.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InclusionObligations.java index 2076a0ea3..e9da21f17 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InclusionObligations.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InclusionObligations.java @@ -137,7 +137,7 @@ public class InclusionObligations implements ISessionStage { Map rgsToUpdate = new HashMap<>(); List execsToUpdate = new ArrayList<>(); boolean loadExecs = !IEnumKey.contains(sessionType, - SessionType.TRDT, SessionType.CURR, SessionType.UNIT, SessionType.IPOT); + SessionType.TRDT, SessionType.CURR, SessionType.PAYM, SessionType.UNIT, SessionType.IPOT); for (Map.Entry> entrySet : registryByGroupId.entrySet()) { log.debug("Processing set of registry with groupId: {}", entrySet.getKey()); String rgsSection = null; @@ -228,6 +228,11 @@ public class InclusionObligations implements ISessionStage { rgsPb.equals("sessionType", SessionType.MEDM.getKey()), rgsPb.equals("sessionType", SessionType.XDEP.getKey()) ); + } else if (SessionType.PAYM.equals(sessionType)) { + return rgsPb.or( + rgsPb.equals("sessionType", SessionType.PREP.getKey()), + rgsPb.equals("sessionType", SessionType.PAYM.getKey()) + ); } else { return rgsPb.equals("sessionType", sessionType.getKey()); } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InspectionObligationsV2.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InspectionObligationsV2.java index 14253e8c8..2bf160ce3 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InspectionObligationsV2.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InspectionObligationsV2.java @@ -163,7 +163,7 @@ public class InspectionObligationsV2 implements ISessionStage { boolean isUncovered = false; List checkResults = new ArrayList<>(); for (Registry obligation : obligationsInGroup) { - if ((SessionType.FINL.equals(sessionType) || SessionType.UNIT.equals(sessionType) || SessionType.MEDM.equals(sessionType)) && RegistryInstrumentType.S.equalsByKey(obligation.getRegistryInstrumentType())) { + if ((SessionType.FINL.equals(sessionType) || SessionType.UNIT.equals(sessionType) || SessionType.MEDM.equals(sessionType)) && RegistryInstrumentType.S.equalsByKey(obligation.getRegistryInstrumentType()) || SessionType.PAYM.equals(sessionType)) { checkResults.add(new CheckResult(obligation, false)); continue; }