clearing-service BalanceRevise ревизии 3 стадии AgainRevise.

This commit is contained in:
AKurakin 2023-09-01 18:11:30 +03:00
parent ef217e8ff7
commit 210b855cb2
13 changed files with 161 additions and 10 deletions

View file

@ -153,8 +153,9 @@ public class StateBnConfig extends EnumStateMachineConfigurerAdapter<TaskType, S
TaskType.InclusionToPool, TaskType.InclusionToPool,
TaskType.InspectionObligations, TaskType.InspectionObligations,
TaskType.FormingRegistersOnOS, TaskType.FormingRegistersOnOS,
TaskType.FormingPaymentInstruction TaskType.FormingPaymentInstruction,
// TaskType.UnlockResources, // TaskType.UnlockResources,
TaskType.AgainRevise
// TaskType.FinishingSession, // TaskType.FinishingSession,
// TaskType.EndStageNotification // TaskType.EndStageNotification
))) )))
@ -202,8 +203,13 @@ public class StateBnConfig extends EnumStateMachineConfigurerAdapter<TaskType, S
.action(formingRegistersOnOSAction) .action(formingRegistersOnOSAction)
.and() .and()
.withExternal() .withExternal()
.source(TaskType.FormingPaymentInstruction).target(TaskType.FormingPaymentInstruction) .source(TaskType.FormingPaymentInstruction).target(TaskType.AgainRevise)
.action(formingPaymentInstructionAction); .action(formingPaymentInstructionAction)
.and()
.withExternal()
.event(SessionEvent.Revise)
.source(TaskType.AgainRevise).target(TaskType.AgainRevise)
.action(balanceReviseAction);
} }

View file

@ -232,6 +232,10 @@ public class FinalMkrSession extends AbstractSession implements InitializingBean
log.error("cannot continue session, current stage is {}", currStage.get()); log.error("cannot continue session, current stage is {}", currStage.get());
throw new StageException(); throw new StageException();
} }
//stage 9 continue revision
{
runStage(TaskType.AgainRevise, currSession.getId(), balanceRevise);
}
//stage 10 //stage 10
{ {
FinishingSessionPayload payload = new FinishingSessionPayload(); FinishingSessionPayload payload = new FinishingSessionPayload();

View file

@ -211,6 +211,10 @@ public class IntermediateMkrSession extends AbstractSession implements Initializ
log.error("cannot continue session, current stage is {}", currStage.get()); log.error("cannot continue session, current stage is {}", currStage.get());
throw new StageException(); throw new StageException();
} }
//stage 9 continue revision
{
runStage(TaskType.AgainRevise, currSession.getId(), balanceRevise);
}
//stage 10 //stage 10
{ {
FinishingSessionPayload payload = new FinishingSessionPayload(); FinishingSessionPayload payload = new FinishingSessionPayload();

View file

@ -203,6 +203,10 @@ public class PrimaryAuctionB0Session extends AbstractSession implements Initiali
log.error("cannot continue session, current stage is {}", currStage.get()); log.error("cannot continue session, current stage is {}", currStage.get());
throw new StageException(); throw new StageException();
} }
//stage 9 continue revision
{
runStage(TaskType.AgainRevise, currSession.getId(), balanceRevise);
}
//stage 10 //stage 10
{ {
FinishingSessionPayload payload = new FinishingSessionPayload(); FinishingSessionPayload payload = new FinishingSessionPayload();

View file

@ -205,6 +205,10 @@ public class PrimaryAuctionBnSession extends AbstractSession implements Initiali
log.error("cannot continue session, current stage is {}", currStage.get()); log.error("cannot continue session, current stage is {}", currStage.get());
throw new StageException(); throw new StageException();
} }
//stage 9 continue revision
{
runStage(TaskType.AgainRevise, currSession.getId(), balanceRevise);
}
//stage 10 //stage 10
{ {
FinishingSessionPayload payload = new FinishingSessionPayload(); FinishingSessionPayload payload = new FinishingSessionPayload();

View file

@ -204,6 +204,10 @@ public class PrimaryAuctionT0Session extends AbstractSession implements Initiali
log.error("cannot continue session, current stage is {}", currStage.get()); log.error("cannot continue session, current stage is {}", currStage.get());
throw new StageException(); throw new StageException();
} }
//stage 9 continue revision
{
runStage(TaskType.AgainRevise, currSession.getId(), balanceRevise);
}
//stage 10 //stage 10
{ {
FinishingSessionPayload payload = new FinishingSessionPayload(); FinishingSessionPayload payload = new FinishingSessionPayload();

View file

@ -183,6 +183,10 @@ public class ReturnDepositSession extends AbstractSession implements Initializin
log.error("cannot continue session, current stage is {}", currStage.get()); log.error("cannot continue session, current stage is {}", currStage.get());
throw new StageException(); throw new StageException();
} }
//stage 9 continue revision
{
runStage(TaskType.AgainRevise, currSession.getId(), balanceRevise);
}
//stage 10 //stage 10
{ {
FinishingSessionPayload payload = new FinishingSessionPayload(); FinishingSessionPayload payload = new FinishingSessionPayload();

View file

@ -199,6 +199,10 @@ public class SecondaryAuctionT0Session extends AbstractSession implements Initia
log.error("cannot continue session, current stage is {}", currStage.get()); log.error("cannot continue session, current stage is {}", currStage.get());
throw new StageException(); throw new StageException();
} }
//stage 9 continue revision
{
runStage(TaskType.AgainRevise, currSession.getId(), balanceRevise);
}
//stage 10 //stage 10
{ {
FinishingSessionPayload payload = new FinishingSessionPayload(); FinishingSessionPayload payload = new FinishingSessionPayload();

View file

@ -41,10 +41,10 @@ public enum TaskType implements IEnumKey {
* step 8 * step 8
*/ */
UnlockResources("CL08"), UnlockResources("CL08"),
// /** /**
// * step 0/9 * step 0/9
// */ */
// AgainRevise("CL09"), AgainRevise("CL09"),
/** /**
* Step 10 * Step 10
*/ */

View file

@ -13,6 +13,7 @@ import ru.clearing.classes.statics.data.statement.Statement;
import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.Consts; 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.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.platform.messaging.service.sender.KafkaSender;
import ru.spcex.clearing.session.stage.ISessionStage; import ru.spcex.clearing.session.stage.ISessionStage;
import ru.spcex.clearing.session.stage.StageResult; import ru.spcex.clearing.session.stage.StageResult;
@ -28,6 +29,9 @@ import java.math.BigDecimal;
import java.time.Instant; import java.time.Instant;
import java.time.temporal.ChronoUnit; import java.time.temporal.ChronoUnit;
import java.util.Collection; import java.util.Collection;
import java.util.Objects;
import static ru.spcex.platform.utils.number.BigDecimalUtil.safeBD;
@Service @Service
@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) @Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE)
@ -60,7 +64,11 @@ public class BalanceRevise implements ISessionStage {
return sendSdfs(); return sendSdfs();
} }
case ContinueRevise -> { case ContinueRevise -> {
return revise(); revise(); // основная сверка
return reviseStage1(); // подготовка к стадии 3 (к AgainRevise)
}
case AgainRevise -> {
return reviseStage3();
} }
default -> throw new IllegalStateException("unknown task " + task.getTaskType()); default -> throw new IllegalStateException("unknown task " + task.getTaskType());
} }
@ -94,8 +102,60 @@ public class BalanceRevise implements ISessionStage {
return new StageResult<>(null, true); return new StageResult<>(null, true);
} }
private BigDecimal safeBD(BigDecimal value) { private StageResult<?> reviseStage1() {
return value != null ? value : BigDecimal.ZERO; String sql = RegistryCodeSqlBuilder.getInstance(ru.spcex.platform.enumeration.RegistryTradingParams.A__T).build();
Collection<Registry> 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<Registry> 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) { private void newSDf56(Statement statement) {

View file

@ -1,5 +1,7 @@
package ru.spcex.clearing.session.stage.impl; 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.Logger;
import org.slf4j.LoggerFactory; import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
@ -31,6 +33,7 @@ import java.util.*;
import java.util.stream.Collectors; import java.util.stream.Collectors;
import static ru.spcex.platform.enumeration.RegistryTradingParams.*; import static ru.spcex.platform.enumeration.RegistryTradingParams.*;
import static ru.spcex.platform.utils.number.BigDecimalUtil.safeBD;
@Service @Service
@ -156,9 +159,52 @@ public class InspectionObligations implements ISessionStage {
} }
} }
stageRevision2(sessionId);
return new StageResult(null, true); 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<Registry> registriesAT = registryImdg.getCollectionObjectsBySQL(sqlAT);
Collection<Registry> registriesAB = registryImdg.getCollectionObjectsBySQL(sqlAB);
log.debug("Select {} registers by \"{}\", {} registers by \"{}\" for revision step 2",
registriesAT.size(), sqlAT, registriesAB.size(), sqlAB);
Map<Triple<Long, Long, Long>, List<Registry>> 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<Long, Long, Long> key = new MutableTriple(regT.getCompanyId(), regT.getAccountId(), regT.getSecurityId());
List<Registry> 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) { private String searchAssetsByObligationSql(Registry obligation) {
RegistryTradingParams counterRegistryTradingParams = null; RegistryTradingParams counterRegistryTradingParams = null;
if (IEnumKey.getEnumByKey(RegistryInstrumentType.class, obligation.getRegistryInstrumentType()) == RegistryInstrumentType.S) { if (IEnumKey.getEnumByKey(RegistryInstrumentType.class, obligation.getRegistryInstrumentType()) == RegistryInstrumentType.S) {

View file

@ -7,6 +7,7 @@ import ru.spcex.platform.enumeration.RegistryDesignation;
import ru.spcex.platform.enumeration.RegistryInstrumentType; import ru.spcex.platform.enumeration.RegistryInstrumentType;
import ru.spcex.platform.enumeration.RegistryUnit; import ru.spcex.platform.enumeration.RegistryUnit;
import java.math.BigDecimal;
import java.util.*; import java.util.*;
public class RegistryUtil { public class RegistryUtil {

View file

@ -24,6 +24,8 @@ public record RegistryTradingParams(RegistryDesignation registryDesignation,
public final static RegistryTradingParams AM_F; public final static RegistryTradingParams AM_F;
public final static RegistryTradingParams AM_T; public final static RegistryTradingParams AM_T;
public final static RegistryTradingParams AM_B; 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 AS_T;
public final static RegistryTradingParams DS_T; public final static RegistryTradingParams DS_T;
public final static RegistryTradingParams AS_B; public final static RegistryTradingParams AS_B;
@ -83,6 +85,14 @@ public record RegistryTradingParams(RegistryDesignation registryDesignation,
RegistryInstrumentType.M, RegistryInstrumentType.M,
null, null,
RegistryUnit.B); 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, AS_T = new RegistryTradingParams(RegistryDesignation.A,
RegistryInstrumentType.S, RegistryInstrumentType.S,
null, null,