From 8955aafdb953605118462488bcf9766223bd860e Mon Sep 17 00:00:00 2001 From: ialbert Date: Thu, 1 Jun 2023 15:44:39 +0300 Subject: [PATCH] ReturnDepositSession --- .../session/stage/ReturnDepositSession.java | 206 +++++++++++++++ .../InspectionObligationsDepositReturn.java | 236 ++++++++++++++++++ .../platform/enumeration/RegistryUnit.java | 1 + 3 files changed, 443 insertions(+) create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/ReturnDepositSession.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InspectionObligationsDepositReturn.java diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/ReturnDepositSession.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/ReturnDepositSession.java new file mode 100644 index 000000000..bd97d12ad --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/ReturnDepositSession.java @@ -0,0 +1,206 @@ +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.misc.Session; +import ru.clearing.classes.statics.data.payment.PaymentInstruction; +import ru.clearing.classes.statics.data.registry.Registry; +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.EndStageNotificationPayload; +import ru.spcex.clearing.session.stage.task.FinishingSessionPayload; +import ru.spcex.clearing.session.stage.task.InclusionToPoolPayload; +import ru.spcex.clearing.session.stage.task.InspectionPoolPayload; +import ru.spcex.platform.enumeration.RegistryStatus; +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; + +@Service +public class ReturnDepositSession 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 InspectionObligationsDepositReturn inspectionObligationsDepositReturn; + private final FormingRegistersOnOS formingRegistersOnOS; + private final FormingPaymentInstruction formingPaymentInstruction; + private final UnlockResources unlockResources; + private final FinishingSession finishingSession; + private final EndStageNotification endStageNotification; + private final Imdg registryImdg; + + + public ReturnDepositSession( + ImdgProvider imdgProvider, + BalanceRevise balanceRevise, + DealsPrepare dealsPrepare, + RequirementsAndObligationCreation requirementsAndObligationCreation, + ObligationAdmission obligationsAdmission, + InclusionObligations inclusionObligations, + FormingRegistersOnOS formingRegistersOnOS, + FormingPaymentInstruction formingPaymentInstruction, + UnlockResources unlockResources, + FinishingSession finishingSession, + EndStageNotification endStageNotification, + IMessageResolver messageResolver, + InspectionObligationsDepositReturn inspectionObligationsDepositReturn) { + 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.inspectionObligationsDepositReturn = inspectionObligationsDepositReturn; + + this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class); + } + + @Override + public void afterPropertiesSet() throws Exception { + ImdgPredicateBuilder rgsPrctBuilder = registryImdg.predicateBuilder(); + inclusionObligations.addRegistryCondition( + rgsPrctBuilder.or(rgsPrctBuilder.equals("registryStatus", RegistryStatus.PROC.getKey()), + rgsPrctBuilder.equals("registryStatus", RegistryStatus.MNG.getKey())) + ); + inclusionObligations.addRegistryCondition( + rgsPrctBuilder.less("valueDate", LocalDate.now()) + ); + + } + + @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 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, inspectionObligationsDepositReturn); + + } + //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(TaskType.StartRevise); + return true; + } + } + } + + @Override + protected Section section() { + return Section.MKR; + } + + @Override + protected SessionType sectionType() { + return SessionType.XDEP; + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InspectionObligationsDepositReturn.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InspectionObligationsDepositReturn.java new file mode 100644 index 000000000..5afa4673f --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InspectionObligationsDepositReturn.java @@ -0,0 +1,236 @@ +package ru.spcex.clearing.session.stage.impl; + +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.registry.Registry; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.session.stage.ISessionStage; +import ru.spcex.clearing.session.stage.StageResult; +import ru.spcex.clearing.session.stage.Task; +import ru.spcex.platform.enumeration.*; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; +import ru.spcex.platform.imdg.api.predicate.specific.RegistryCodeSqlBuilder; +import ru.spcex.platform.utils.enumeration.IEnumKey; + +import java.math.BigDecimal; +import java.time.Instant; +import java.util.List; +import java.util.Map; +import java.util.Optional; +import java.util.stream.Collectors; + +import static ru.spcex.platform.utils.number.BigDecimalUtil.safeBD; + + +@Service +public class InspectionObligationsDepositReturn implements ISessionStage { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final Imdg registryImdg; + + private final static RegistryTradingParams OS_T; + private final static RegistryTradingParams OM_T; + private final static RegistryTradingParams TS_T; + private final static RegistryTradingParams TM_T; + private final static RegistryTradingParams DM_X; + private final static RegistryTradingParams DM_T; + private final static RegistryTradingParams AM_F; + + static { + OS_T = new RegistryTradingParams(RegistryDesignation.O, + RegistryInstrumentType.S, + null, + RegistryUnit.T); + + OM_T = new RegistryTradingParams(RegistryDesignation.O, + RegistryInstrumentType.M, + null, + RegistryUnit.T); + + TS_T = new RegistryTradingParams(RegistryDesignation.T, + RegistryInstrumentType.S, + null, + RegistryUnit.T); + + TM_T = new RegistryTradingParams(RegistryDesignation.T, + RegistryInstrumentType.M, + null, + RegistryUnit.T); + DM_X = new RegistryTradingParams(RegistryDesignation.D, + RegistryInstrumentType.M, + null, + RegistryUnit.X); + DM_T = new RegistryTradingParams(RegistryDesignation.D, + RegistryInstrumentType.M, + null, + RegistryUnit.T); + AM_F = new RegistryTradingParams(RegistryDesignation.A, + RegistryInstrumentType.M, + null, + RegistryUnit.F); + + + } + + @Autowired + public InspectionObligationsDepositReturn(ImdgProvider imdgProvider) { + this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class); + } + + @Override + public StageResult submit(Task task) { + switch (task.getTaskType()) { + case InspectionObligations -> { + return inspectionObligations(); + } + default -> throw new IllegalStateException("Unknown task type: " + task.getTaskType()); + } + } + + private StageResult inspectionObligations() { + String sqlCondition = String.format("(%s) and registryStatus = '%s'", + RegistryCodeSqlBuilder.getInstance(OS_T, OM_T, TS_T, TM_T).build(), + RegistryStatus.POOL.getKey()); + + Map> registriesByGroup = registryImdg.getCollectionObjectsBySQL(sqlCondition) + .stream(). + collect(Collectors.groupingBy(Registry::getGroupId)); + log.info("registry groups found {}", registriesByGroup.size()); + + for (Map.Entry> entry : registriesByGroup.entrySet()) { + List group = entry.getValue(); + Optional omtInGroupO = group.stream().filter(registry -> equalByRgs(OM_T, registry)).findFirst(); + if (omtInGroupO.isEmpty()) { + log.error("groupId {} failed to find OM*T registry", entry.getKey()); + group.forEach(rgs -> updateStatus(rgs, RegistryStatus.MNG)); + continue; + } + Registry omtRgs = omtInGroupO.get(); + Optional amfAssetO = searchAssetByOMT(omtRgs); + if (amfAssetO.isEmpty()) { + log.error("OM*T register.id={} groupId={} failed to find AM*F asset", omtRgs.getId(), entry.getKey()); + group.forEach(rgs -> updateStatus(rgs, RegistryStatus.MNG)); + continue; + } + Registry amfAsset = amfAssetO.get(); + log.debug("groupId={}, OM*T.id={}, AM*F.id={}", entry.getKey(), omtRgs.getId(), amfAsset.getId()); + BigDecimal omtBalance = safeBD(omtRgs.getBalance()); + BigDecimal amfBalance = safeBD(amfAsset.getBalance()); + + Optional dmx = searchDmx(omtRgs); + if (dmx.isPresent()) { + BigDecimal dmxBalance = safeBD(dmx.get().getBalance()); + log.debug("OM*T.id={} -> DM*X.id={}, dmxBalance={}, omtBalance={}, amfBalance={}", + omtRgs.getId(), + dmx.get().getId(), + dmxBalance, omtBalance, amfBalance + ); + if (dmxBalance.compareTo(omtBalance) >= 0 && amfBalance.compareTo(omtBalance) >= 0) { + group.forEach(rgs -> updateStatus(rgs, RegistryStatus.OK)); + } else { + group.forEach(rgs -> updateStatus(rgs, RegistryStatus.MNG)); + } + continue; + } + Optional dmtInfo = searchDmtInfo(omtRgs); + if (dmtInfo.isPresent()) { + BigDecimal dmtBalance = safeBD(dmtInfo.get().getBalance()); + log.debug("OM*T.id={} -> DM*T(INFO).id={}, dmtBalance={}, omtBalance={}, amfBalance={}", + omtRgs.getId(), + dmtInfo.get().getId(), + dmtBalance, omtBalance, amfBalance + ); + if (dmtBalance.compareTo(omtBalance) >= 0 && amfBalance.compareTo(omtBalance) >=0) { + group.forEach(rgs -> updateStatus(rgs, RegistryStatus.OK)); + } else { + group.forEach(rgs -> updateStatus(rgs, RegistryStatus.MNG)); + } + continue; + } + Optional dmtClnr = searchDmtClrn(omtRgs); + if (dmtClnr.isPresent()) { + BigDecimal dmtBalance = safeBD(dmtClnr.get().getBalance()); + log.debug("OM*T.id={} -> DM*T(CLNR).id={}, dmtBalance={}, omtBalance={}, amfBalance={}", + omtRgs.getId(), + dmtClnr.get().getId(), + dmtBalance, omtBalance, amfBalance + ); + if (dmtBalance.compareTo(omtBalance) >= 0 && amfBalance.compareTo(omtBalance) >= 0) { + group.forEach(rgs -> updateStatus(rgs, RegistryStatus.OK)); + } else { + group.forEach(rgs -> updateStatus(rgs, RegistryStatus.MNG)); + } + continue; + } + log.info("OM*T.id={} -> no DM*X/DM*T(INFO/CLRN) registry found. setting MNG status to group", omtRgs.getId()); + group.forEach(rgs -> updateStatus(rgs, RegistryStatus.MNG)); + } + return new StageResult<>(null, true); + } + + private void updateStatus(Registry registry, RegistryStatus registryStatus) { + log.trace("Update registry.id: {} to {}", registry.getId(), registryStatus.getKey()); + registry.setRegistryStatus(registryStatus.getKey()); + registry.setUpdated(Instant.now()); + registryImdg.update(registry); + } + + /** + * Если УК зачисляет средства на ТБС Инициатора + */ + private Optional searchDmx(Registry rgs) { + String sqlCondition = String.format("(%s) and companyId = %d", + RegistryCodeSqlBuilder.getInstance(DM_X).build(), + rgs.getCounterPartyId()); + Registry dmx = registryImdg.getSingleObjectBySQL(sqlCondition); + return Optional.ofNullable(dmx); + } + + /** + * Если УК зачисляет средства на свой регистр на КС + */ + private Optional searchDmtInfo(Registry rgs) { + String sqlCondition = String.format("(%s) and accountType='%s' and companyId = %d and counterPartyId = %d", + RegistryCodeSqlBuilder.getInstance(DM_T).build(), + AccountType.Info.getKey(), + rgs.getCompanyId(), + rgs.getCounterPartyId()); + Registry dmt = registryImdg.getSingleObjectBySQL(sqlCondition); + return Optional.ofNullable(dmt); + } + + /** + * Если УК зачисляет средства на свой ТБС + */ + private Optional searchDmtClrn(Registry rgs) { + String sqlCondition = String.format("(%s) and accountType='%s' and companyId = %d and counterPartyId = %d", + RegistryCodeSqlBuilder.getInstance(DM_T).build(), + AccountType.Clrn.getKey(), + rgs.getCompanyId(), + rgs.getCounterPartyId()); + Registry dmt = registryImdg.getSingleObjectBySQL(sqlCondition); + return Optional.ofNullable(dmt); + } + + private Optional searchAssetByOMT(Registry obligation) { + String sqlCondition = String.format("%s and " + + "tradingClearingRegistryId = '%s' and " + + "companyId = '%s'", + RegistryCodeSqlBuilder.getInstance(AM_F).build(), + obligation.getTradingClearingRegistryId(), + obligation.getCompanyId()); + Registry amf = registryImdg.getSingleObjectBySQL(sqlCondition); + return Optional.ofNullable(amf); + } + + private static boolean equalByRgs(RegistryTradingParams rgsParams, Registry rgs) { + return rgsParams.equalByRegistry( + IEnumKey.getEnumByKey(RegistryDesignation.class, rgs.getRegistryDesignation()), + IEnumKey.getEnumByKey(RegistryInstrumentType.class, rgs.getRegistryInstrumentType()), + IEnumKey.getEnumByKey(RegistryCapacity.class, rgs.getRegistryCapacity()), + IEnumKey.getEnumByKey(RegistryUnit.class, rgs.getRegistryUnit()) + ); + } +} diff --git a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/RegistryUnit.java b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/RegistryUnit.java index 9fe5baa20..230107b42 100644 --- a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/RegistryUnit.java +++ b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/RegistryUnit.java @@ -8,6 +8,7 @@ public enum RegistryUnit implements IEnumKey { R("R"), F("F"), B("B"), + X("X"), ; private final String key;