diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/session/StateBnConfig.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/session/StateBnConfig.java index f5c7ef729..7f685f721 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/session/StateBnConfig.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/session/StateBnConfig.java @@ -152,14 +152,16 @@ public class StateBnConfig extends EnumStateMachineConfigurerAdapter(Arrays.asList(TaskType.StartRevise, TaskType.ContinueRevise, + TaskType.StartRevisePart1, TaskType.DealsPrepare, TaskType.RequirementsAndObligationsCreate, TaskType.ObligationsAdmission, TaskType.InclusionToPool, TaskType.InspectionObligations, TaskType.FormingRegistersOnOS, - TaskType.FormingPaymentInstruction + TaskType.FormingPaymentInstruction, // TaskType.UnlockResources, + TaskType.AgainRevise // TaskType.FinishingSession, // TaskType.EndStageNotification ))) @@ -177,11 +179,16 @@ public class StateBnConfig extends EnumStateMachineConfigurerAdapter { log.info("SDF57 and SDF01 received, continue session"); }) .and() + .withExternal() + .event(SessionEvent.Revise) + .source(TaskType.StartRevisePart1).target(TaskType.DealsPrepare) + .action(balanceReviseAction) + .and() .withExternal() .source(TaskType.DealsPrepare).target(TaskType.RequirementsAndObligationsCreate) .action(dealPrepareAction) @@ -207,8 +214,13 @@ public class StateBnConfig extends EnumStateMachineConfigurerAdapter> dealsPreparationResult; { @@ -232,6 +233,10 @@ public class FinalMkrSession extends AbstractSession implements InitializingBean log.error("cannot continue session, current stage is {}", currStage.get()); throw new StageException(); } + //stage 9 continue revision + { + runStage(TaskType.AgainRevise, currSession.getId(), balanceRevise); + } //stage 10 { FinishingSessionPayload payload = new FinishingSessionPayload(); 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 988afa6f6..e22e148ff 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 @@ -153,6 +153,7 @@ public class IntermediateMkrSession extends AbstractSession implements Initializ log.error("cannot continue session, current stage is {}", currStage.get()); throw new StageException(); } + runStage(TaskType.StartRevisePart1, currSession.getId(), balanceRevise); //stage 1 StageResult> dealsPreparationResult; { @@ -211,6 +212,10 @@ public class IntermediateMkrSession extends AbstractSession implements Initializ log.error("cannot continue session, current stage is {}", currStage.get()); throw new StageException(); } + //stage 9 continue revision + { + runStage(TaskType.AgainRevise, currSession.getId(), balanceRevise); + } //stage 10 { FinishingSessionPayload payload = new FinishingSessionPayload(); diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/PrimaryAuctionB0Session.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/PrimaryAuctionB0Session.java index bac39ecb9..2b7a51980 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/PrimaryAuctionB0Session.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/PrimaryAuctionB0Session.java @@ -147,6 +147,7 @@ public class PrimaryAuctionB0Session extends AbstractSession implements Initiali log.error("cannot continue session, current stage is {}", currStage.get()); throw new StageException(); } + runStage(TaskType.StartRevisePart1, currSession.getId(), balanceRevise); //stage 1 StageResult> dealsPreparationResult; { @@ -203,6 +204,10 @@ public class PrimaryAuctionB0Session extends AbstractSession implements Initiali log.error("cannot continue session, current stage is {}", currStage.get()); throw new StageException(); } + //stage 9 continue revision + { + runStage(TaskType.AgainRevise, currSession.getId(), balanceRevise); + } //stage 10 { FinishingSessionPayload payload = new FinishingSessionPayload(); 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 1dcb36d75..ccfc80613 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 @@ -148,6 +148,7 @@ public class PrimaryAuctionBnSession extends AbstractSession implements Initiali log.error("cannot continue session, current stage is {}", currStage.get()); throw new StageException(); } + runStage(TaskType.StartRevisePart1, currSession.getId(), balanceRevise); //stage 1 StageResult> dealsPreparationResult; { @@ -205,6 +206,10 @@ public class PrimaryAuctionBnSession extends AbstractSession implements Initiali log.error("cannot continue session, current stage is {}", currStage.get()); throw new StageException(); } + //stage 9 continue revision + { + runStage(TaskType.AgainRevise, currSession.getId(), balanceRevise); + } //stage 10 { FinishingSessionPayload payload = new FinishingSessionPayload(); diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/PrimaryAuctionT0Session.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/PrimaryAuctionT0Session.java index 8261e91c9..df42c6510 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/PrimaryAuctionT0Session.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/PrimaryAuctionT0Session.java @@ -148,6 +148,7 @@ public class PrimaryAuctionT0Session extends AbstractSession implements Initiali log.error("cannot continue session, current stage is {}", currStage.get()); throw new StageException(); } + runStage(TaskType.StartRevisePart1, currSession.getId(), balanceRevise); //stage 1 StageResult> dealsPreparationResult; { @@ -204,6 +205,10 @@ public class PrimaryAuctionT0Session extends AbstractSession implements Initiali log.error("cannot continue session, current stage is {}", currStage.get()); throw new StageException(); } + //stage 9 continue revision + { + runStage(TaskType.AgainRevise, currSession.getId(), balanceRevise); + } //stage 10 { FinishingSessionPayload payload = new FinishingSessionPayload(); 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 index aa433d809..b808cc95a 100644 --- 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 @@ -139,6 +139,7 @@ public class ReturnDepositSession extends AbstractSession implements Initializin log.error("cannot continue session, current stage is {}", currStage.get()); throw new StageException(); } + runStage(TaskType.StartRevisePart1, currSession.getId(), balanceRevise); //stage 4 { InclusionToPoolPayload inclusionToPoolPayload = new InclusionToPoolPayload(); @@ -183,6 +184,10 @@ public class ReturnDepositSession extends AbstractSession implements Initializin log.error("cannot continue session, current stage is {}", currStage.get()); throw new StageException(); } + //stage 9 continue revision + { + runStage(TaskType.AgainRevise, currSession.getId(), balanceRevise); + } //stage 10 { FinishingSessionPayload payload = new FinishingSessionPayload(); 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 index 624b63913..d6a8923ad 100644 --- 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 @@ -142,6 +142,7 @@ public class SecondaryAuctionT0Session extends AbstractSession implements Initia log.error("cannot continue session, current stage is {}", currStage.get()); throw new StageException(); } + runStage(TaskType.StartRevisePart1, currSession.getId(), balanceRevise); //stage 1 StageResult> dealsPreparationResult; { @@ -199,6 +200,10 @@ public class SecondaryAuctionT0Session extends AbstractSession implements Initia log.error("cannot continue session, current stage is {}", currStage.get()); throw new StageException(); } + //stage 9 continue revision + { + runStage(TaskType.AgainRevise, currSession.getId(), balanceRevise); + } //stage 10 { FinishingSessionPayload payload = new FinishingSessionPayload(); diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/TaskType.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/TaskType.java index 9353c085d..142bca615 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/TaskType.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/TaskType.java @@ -8,6 +8,7 @@ public enum TaskType implements IEnumKey { */ StartRevise("CLR0"), ContinueRevise("CLR 0_0"), + StartRevisePart1("CLR 0_9"), /** * step 1 @@ -41,10 +42,10 @@ public enum TaskType implements IEnumKey { * step 8 */ UnlockResources("CL08"), -// /** -// * step 0/9 -// */ -// AgainRevise("CL09"), + /** + * step 0/9 + */ + AgainRevise("CL09"), /** * Step 10 */ diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/BalanceRevise.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/BalanceRevise.java index fd9738c2f..c2a17fe40 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/BalanceRevise.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/BalanceRevise.java @@ -13,6 +13,7 @@ import ru.clearing.classes.statics.data.statement.Statement; import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.platform.messaging.domain.Consts; import ru.spcex.clearing.platform.messaging.domain.cud.clearing.SdfClearingRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.utilities.NotificationNewRequest; import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.clearing.session.stage.ISessionStage; import ru.spcex.clearing.session.stage.StageResult; @@ -28,6 +29,9 @@ import java.math.BigDecimal; import java.time.Instant; import java.time.temporal.ChronoUnit; import java.util.Collection; +import java.util.Objects; + +import static ru.spcex.platform.utils.number.BigDecimalUtil.safeBD; @Service @Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) @@ -60,7 +64,13 @@ public class BalanceRevise implements ISessionStage { return sendSdfs(); } case ContinueRevise -> { - return revise(); + return revise(); // основная сверка + } + case StartRevisePart1 -> { // часть 1, часть 2 на стадии 5 (InspectionObligations) + return reviseStage1(); // подготовка к стадии 3 (к AgainRevise) + } + case AgainRevise -> { // часть 3 + return reviseStage3(); } default -> throw new IllegalStateException("unknown task " + task.getTaskType()); } @@ -94,8 +104,60 @@ public class BalanceRevise implements ISessionStage { return new StageResult<>(null, true); } - private BigDecimal safeBD(BigDecimal value) { - return value != null ? value : BigDecimal.ZERO; + private StageResult reviseStage1() { + String sql = RegistryCodeSqlBuilder.getInstance(ru.spcex.platform.enumeration.RegistryTradingParams.A__T).build(); + Collection regsAT = registryImdg.getCollectionObjectsBySQL(sql); + log.trace("Select {} registry's by query \"{}\" for revision step 1", regsAT.size(), sql); + + Instant now = Instant.now(); + int updateCount = 0; + for (Registry reg : regsAT) { + if (reg.getBalance() == null) { + log.debug("Registry[{}] with null balance", reg.getId()); + } else { + if (!Objects.equals(reg.getPlanBalance(), reg.getBalance())) { + reg.setPlanBalance(reg.getBalance()); + reg.setUpdated(now); + registryImdg.update(reg); + updateCount++; + } + } + } + log.debug("At revision stage 1 do updated {} of {} registers {}", updateCount, regsAT.size(), sql); + return new StageResult<>(null, true); + } + + private StageResult reviseStage3() { + String sql = RegistryCodeSqlBuilder.getInstance(RegistryTradingParams.A__T).build(); + Collection regsAT = registryImdg.getCollectionObjectsBySQL(sql); + log.trace("Select {} registry's by query \"{}\" for revision step 3", regsAT.size(), sql); + + int errorRegs = 0; + for (Registry reg : regsAT) { + if (reg.getBalance() == null && reg.getPlanBalance() == null) { + log.debug("Can not verify registry[{}] with null balance and planBalance", reg.getId()); + } else { + if (reg.getBalance() == null || reg.getPlanBalance() == null) { + log.warn("Can not verify registry[{}] with null balance xor planBalance", reg.getId()); + } + if (safeBD(reg.getBalance()).compareTo(safeBD(reg.getPlanBalance())) != 0) { + log.info("Revision: registry id={}, companyId={}, balance={}, plannedBalance={}", + reg.getId(), reg.getCompanyId(), reg.getBalance(), reg.getPlanBalance()); + errorRegs++; + } + } + + } + if (errorRegs > 0) { + log.warn("После сверки обнаружена разница между плановым и фактическим балансом. Всего {} регистров не совпали.", errorRegs); + NotificationNewRequest nRequest = new NotificationNewRequest(); + nRequest.setObjectType(ObjectType.rgst.getKey()); + nRequest.setPriority(Priority.LOW.getKey()); + nRequest.setComment("После сверки обнаружена разница между плановым и фактическим балансом"); + kafkaSender.sendRequestToQueue(Consts.NOTIFICATION_NEW, nRequest); + } + + return new StageResult<>(null, true); } private void newSDf56(Statement statement) { diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InspectionObligations.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InspectionObligations.java index c83bc9767..46c77388f 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InspectionObligations.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InspectionObligations.java @@ -1,5 +1,7 @@ package ru.spcex.clearing.session.stage.impl; +import org.apache.commons.lang3.tuple.MutableTriple; +import org.apache.commons.lang3.tuple.Triple; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Autowired; @@ -31,6 +33,7 @@ import java.util.*; import java.util.stream.Collectors; import static ru.spcex.platform.enumeration.RegistryTradingParams.*; +import static ru.spcex.platform.utils.number.BigDecimalUtil.safeBD; @Service @@ -156,9 +159,52 @@ public class InspectionObligations implements ISessionStage { } } + stageRevision2(sessionId); + return new StageResult(null, true); } + protected void stageRevision2(Long sessionId) { + /* todo logic + Collection rgsAT = findRegistryATBySession // registry_disignation=A and registry_unit=T +Collection rgsAB = findRegistryATBySession // registry_disignation=A and registry_unit=B and sessionId = current + +foreach( at: rgsAT) + foreach (ab: rgsAB) + if (at.companyId=ab.companyId and at.accountId=ab.accountId and at.securityId=ab.securityId) + at.plannecBallance -= ab.balance + */ + String sqlAT = RegistryCodeSqlBuilder.getInstance(A__T).build(); + String sqlAB = String.format("(%s) and sessionId = %d", + RegistryCodeSqlBuilder.getInstance(A__B).build(), + sessionId + ); + Collection registriesAT = registryImdg.getCollectionObjectsBySQL(sqlAT); + Collection registriesAB = registryImdg.getCollectionObjectsBySQL(sqlAB); + log.debug("Select {} registers by \"{}\", {} registers by \"{}\" for revision step 2", + registriesAT.size(), sqlAT, registriesAB.size(), sqlAB); + Map, List> regABIndex = registriesAB.stream().collect(Collectors.groupingBy( + (Registry reg) -> new MutableTriple(reg.getCompanyId(), reg.getAccountId(), reg.getSecurityId()) + )); + int updateCount = 0; + Instant now = Instant.now(); + for (Registry regT : registriesAT) { + Triple key = new MutableTriple(regT.getCompanyId(), regT.getAccountId(), regT.getSecurityId()); + List regsB = regABIndex.get(key); + if (regsB == null) { + log.debug("Registry A__B for registry[{}] (A__T key {}) not found", regT, key); + } else { + for (Registry regB : regsB) { + regT.setPlanBalance(safeBD(regT.getPlanBalance()).subtract(safeBD(regB.getBalance()))); + } + regT.setUpdated(now); + registryImdg.update(regT); + updateCount++; + } + } + log.debug("Updated {} registers A__T with planBalance at {}", updateCount, now); + } + private String searchAssetsByObligationSql(Registry obligation) { RegistryTradingParams counterRegistryTradingParams = null; if (IEnumKey.getEnumByKey(RegistryInstrumentType.class, obligation.getRegistryInstrumentType()) == RegistryInstrumentType.S) { diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/util/RegistryUtil.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/util/RegistryUtil.java index 4361b6d4b..ab1826862 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/util/RegistryUtil.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/util/RegistryUtil.java @@ -7,6 +7,7 @@ import ru.spcex.platform.enumeration.RegistryDesignation; import ru.spcex.platform.enumeration.RegistryInstrumentType; import ru.spcex.platform.enumeration.RegistryUnit; +import java.math.BigDecimal; import java.util.*; public class RegistryUtil { diff --git a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/RegistryTradingParams.java b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/RegistryTradingParams.java index ed86625cd..ecd362ce2 100644 --- a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/RegistryTradingParams.java +++ b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/RegistryTradingParams.java @@ -24,6 +24,8 @@ public record RegistryTradingParams(RegistryDesignation registryDesignation, public final static RegistryTradingParams AM_F; public final static RegistryTradingParams AM_T; public final static RegistryTradingParams AM_B; + public final static RegistryTradingParams A__B; + public final static RegistryTradingParams A__T; public final static RegistryTradingParams AS_T; public final static RegistryTradingParams DS_T; public final static RegistryTradingParams AS_B; @@ -83,6 +85,14 @@ public record RegistryTradingParams(RegistryDesignation registryDesignation, RegistryInstrumentType.M, null, RegistryUnit.B); + A__B = new RegistryTradingParams(RegistryDesignation.A, + null, + null, + RegistryUnit.B); + A__T = new RegistryTradingParams(RegistryDesignation.A, + null, + null, + RegistryUnit.T); AS_T = new RegistryTradingParams(RegistryDesignation.A, RegistryInstrumentType.S, null,