From e88eef031c8f133671f7a5f6f88bd5eb65233931 Mon Sep 17 00:00:00 2001 From: ialbert Date: Fri, 3 Nov 2023 18:57:41 +0300 Subject: [PATCH] kill session command --- .../clearing/config/ValidationConfig.java | 14 ++ .../clearing/service/EventsReceiver.java | 9 +- .../SessionTerminationValidationRule.java | 55 +++++++ .../service/validation/ValidationStored.java | 2 + .../session/stage/AbstractSession.java | 5 + .../session/stage/SessionManager.java | 38 +++-- .../session/stage/SessionTerminator.java | 152 ++++++++++++++++++ .../enumeration/RegistryTradingParams.java | 15 ++ .../platform/messaging/domain/Consts.java | 1 + 9 files changed, 279 insertions(+), 12 deletions(-) create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/SessionTerminationValidationRule.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/SessionTerminator.java diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/ValidationConfig.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/ValidationConfig.java index 053250918..c1f8999b2 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/ValidationConfig.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/ValidationConfig.java @@ -46,6 +46,7 @@ import ru.spcex.platform.utils.validation.ValidatorImpl; import java.util.function.BiFunction; import java.util.function.Function; +import java.util.function.Supplier; @Configuration public class ValidationConfig { @@ -490,6 +491,19 @@ public class ValidationConfig { }; } + @Bean("terminateSessionValidator") + public Supplier terminateSessionValidator() { + return () -> { + ImdgValidationContext ctx = new ImdgValidationContext<>(); + ctx.setLogPrefix(LogPrefixId.INSTANCE); + ctx.addImdg(IMDGDistributedNames.Map_Session, imdgSession); + return new ValidatorImpl<>(ctx, + SessionTerminationValidationRule.SessionPresent, + SessionTerminationValidationRule.StatusCheck + ); + }; + } + @Bean("userRoleVerification") public UserRoleVerification userRoleVerification(ImdgProvider imdgProvider, IMessageResolver msgs) { return new UserRoleVerification(imdgProvider, msgs, ClearingError.UserVerifyDenial); diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/EventsReceiver.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/EventsReceiver.java index 31f8818ed..4527e3f19 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/EventsReceiver.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/EventsReceiver.java @@ -54,6 +54,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean { private final BalanceRevise balanceRevise; private final Sdf05Sender sdf05Sender; private final StatementServiceV2 statementService; + private final SessionTerminator sessionTerminator; private final PaymentInstructionOutboundService pmtOutboundService; @Autowired @@ -65,7 +66,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean { SecondaryAuctionT0Session secondaryAuctionT0Session, PrimaryAuctionB0Session primaryAuctionB0Session, PrimaryAuctionT0Session primaryAuctionT0Session, IntermediateMkrSession intermediateMkrSession, FinalMkrSession finalMkrSession, ReturnDepositSession returnDepositSession, SessionManager sessionManager, Sdf06Executor sdf06Executor, - Sdf10Executor sdf10Executor, BalanceRevise balanceRevise, Sdf05Sender sdf05Sender, StatementServiceV2 statementService, PaymentInstructionOutboundService pmtOutboundService) { + Sdf10Executor sdf10Executor, BalanceRevise balanceRevise, Sdf05Sender sdf05Sender, StatementServiceV2 statementService, SessionTerminator sessionTerminator, PaymentInstructionOutboundService pmtOutboundService) { super(kafkaQueue, kafkaResponseQueue); this.errorResolver = errorResolver; this.clearingService = clearingService; @@ -83,6 +84,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean { this.balanceRevise = balanceRevise; this.sdf05Sender = sdf05Sender; this.statementService = statementService; + this.sessionTerminator = sessionTerminator; this.pmtOutboundService = pmtOutboundService; } @@ -184,6 +186,11 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean { callback(LauncherCommandRequest.class) .setConsumer(task -> sdf05Sender.sendSdf05("9")) .forDestination(Task.sdf05WithCode9Final.topic(), callbacks::put); + + callback(Object.class) + .setFunction(sessionTerminator::stopCurrentSession) + .forDestination(Consts.KILL_SESSION, callbacks::put); + init(); } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/SessionTerminationValidationRule.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/SessionTerminationValidationRule.java new file mode 100644 index 000000000..2f42247ec --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/SessionTerminationValidationRule.java @@ -0,0 +1,55 @@ +package ru.spcex.clearing.service.validation; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import ru.clearing.classes.statics.data.misc.Session; +import ru.spcex.clearing.error.ClearingError; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.session.stage.TaskType; +import ru.spcex.platform.enumeration.SessionStatus; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.validation.ImdgValidationContext; +import ru.spcex.platform.utils.enumeration.EnumMessage; +import ru.spcex.platform.utils.enumeration.IEnumKey; +import ru.spcex.platform.utils.validation.IValidationRule; + +import java.util.Map; +import java.util.Optional; + +public enum SessionTerminationValidationRule implements IValidationRule> { + SessionPresent() { + @Override + public Optional validate(ImdgValidationContext context) { + Imdg sessionImdg = context.obtainMap(IMDGDistributedNames.Map_Session, Session.class); + Session session = sessionImdg.getFirstObjectByFieldValues(Map.of( + "workflowStatus", SessionStatus.ACTV.getKey() + )); + if (session == null) { + log.warn("Cannot stop session, failed to find workflowStatus 'ACTV'"); + return of(ClearingError.RecordNotFound, "active session"); + } + context.storeObject(ValidationStored.SessionTerminate, session); + context.setValidatedObject(session); + return empty(); + } + }, + StatusCheck() { + @Override + public Optional validate(ImdgValidationContext context) { + Session session = context.getStoredObject(ValidationStored.SessionTerminate); + TaskType status = IEnumKey.getEnumByKey(TaskType.class, session.getSessionStatus()); + if (status == null || status.ordinal() >= TaskType.FormingPaymentInstruction.ordinal()) { + log.warn("Cannot stop session.id={}, unknown sessionStatus '{}'", + session.getId(), session.getSessionStatus()); + return of(ClearingError.IncorrectValue, "sessionStatus"); + } + return empty(); + } + } + ; + private final static Logger log = LoggerFactory.getLogger(SessionTerminationValidationRule.class); + @Override + public String ruleName() { + return "SessionTerminationValidationRule." + name(); + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/ValidationStored.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/ValidationStored.java index 115f28e5c..2977937c5 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/ValidationStored.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/ValidationStored.java @@ -32,5 +32,7 @@ public enum ValidationStored { SecurityBySecurityCode, + SessionTerminate, + RegistrysByContract, SplitDepositMaxNumber } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/AbstractSession.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/AbstractSession.java index ae8f646b8..7f06880fb 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/AbstractSession.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/AbstractSession.java @@ -105,6 +105,11 @@ public abstract class AbstractSession { protected void continueRunning(TaskType t) { synchronized (this.currStage) { + //if (WorkflowStatus.Blocked.equalsByKey(this.currSession.getWorkflowStatus())) { + // log.error("cannot run task {}.{}. session.id={} workflowStatus is {}", + // t, t.getKey(), this.currSession.getId(), this.currSession.getWorkflowStatus()); + // throw new StageException(); + //} this.currStage.set(t); this.currSession.setSessionStatus(t.getKey()); this.currSession.setUpdated(Instant.now()); diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/SessionManager.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/SessionManager.java index e55b0c9cf..5294808d8 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/SessionManager.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/SessionManager.java @@ -66,17 +66,7 @@ public class SessionManager { return; } BaseRequest baseRequest = new BaseRequest<>(); - AbstractSession session = null; - switch (sessionType){ - case IPOB -> session = primaryAuctionBnSession; - case IPOT -> session = primaryAuctionT0Session; - case IPO0 -> session = primaryAuctionB0Session; - case TRDT -> session = secondaryAuctionT0Session; - case MEDM -> session = intermediateMkrSession; - case FINL -> session = finalMkrSession; - case XDEP -> session = returnDepositSession; - } - + AbstractSession session = sessionByType(sessionType); if (session != null) { //checkActive checkAllowSessionStart(); @@ -99,4 +89,30 @@ public class SessionManager { throw new ValidationException(err); } } + + public void endSession(Session session) { + SessionType sessionType = IEnumKey.getEnumByKey(SessionType.class, session.getSessionType()); + String err = "failed to stop session.id=%d: unknown sessionType: %s" + .formatted(session.getId(), session.getSessionType()); + if (sessionType == null) { + throw new IllegalStateException(err); + } + AbstractSession sessionFlow = sessionByType(sessionType); + if (sessionFlow == null) throw new IllegalStateException(err); + sessionFlow.endSession(); + } + + private AbstractSession sessionByType(SessionType sessionType) { + AbstractSession session = null; + switch (sessionType){ + case IPOB -> session = primaryAuctionBnSession; + case IPOT -> session = primaryAuctionT0Session; + case IPO0 -> session = primaryAuctionB0Session; + case TRDT -> session = secondaryAuctionT0Session; + case MEDM -> session = intermediateMkrSession; + case FINL -> session = finalMkrSession; + case XDEP -> session = returnDepositSession; + } + return session; + } } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/SessionTerminator.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/SessionTerminator.java new file mode 100644 index 000000000..8f1ec54a1 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/SessionTerminator.java @@ -0,0 +1,152 @@ +package ru.spcex.clearing.session.stage; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +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.ExecutionDeposit; +import ru.clearing.classes.statics.data.execution.ExecutionFond; +import ru.clearing.classes.statics.data.misc.Session; +import ru.clearing.classes.statics.data.registry.Registry; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.notification.NotificationSender; +import ru.spcex.clearing.platform.messaging.domain.BaseRequest; +import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate; +import ru.spcex.clearing.service.registry.AssetTBFProcessing; +import ru.spcex.clearing.service.validation.ValidationStored; +import ru.spcex.clearing.util.services.RequestHelper; +import ru.spcex.platform.classes.base.SpcexObjectBase; +import ru.spcex.platform.enumeration.ObjectType; +import ru.spcex.platform.enumeration.Priority; +import ru.spcex.platform.enumeration.RegistryTradingParams; +import ru.spcex.platform.enumeration.Section; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; +import ru.spcex.platform.imdg.api.predicate.ImdgPredicate; +import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder; +import ru.spcex.platform.imdg.api.predicate.specific.RegistryCodeSqlBuilder; +import ru.spcex.platform.utils.enumeration.EnumMessage; +import ru.spcex.platform.utils.enumeration.IMessageResolver; +import ru.spcex.platform.utils.validation.IValidator; + +import java.math.BigDecimal; +import java.time.Instant; +import java.util.Collection; +import java.util.Optional; +import java.util.function.Consumer; +import java.util.function.Supplier; + +@Service +public class SessionTerminator { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final Imdg execDpstImdg; + private final Imdg execFondImdg; + private final Imdg rgsImdg; + + private final NotificationSender notification; + private final AssetTBFProcessing assets; + private final IMessageResolver msgs; + private final RequestHelper reqInfo; + private final Supplier validation; + private final SessionManager sessionMng; + + + public SessionTerminator(ImdgProvider imdgProvider, + NotificationSender notification, + AssetTBFProcessing assets, IMessageResolver msgs, + RequestHelper reqInfo, + @Qualifier("terminateSessionValidator") Supplier validation, SessionManager sessionMng) { + this.notification = notification; + this.assets = assets; + this.msgs = msgs; + this.execDpstImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionDeposit, ExecutionDeposit.class); + this.execFondImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionFond, ExecutionFond.class); + this.rgsImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class); + this.reqInfo = reqInfo; + this.validation = validation; + this.sessionMng = sessionMng; + } + + public RequestInfoUpdate stopCurrentSession(BaseRequest req) { + IValidator validator = validation.get(); + Optional err = validator.tillFirstError(); + if (err.isPresent()) { + //fixme - уточнить нужно ли + notification.sendNotification(ObjectType.session, msgs.resolve(err.get()), Priority.HIGH); + return reqInfo.error(req.getId(), err.get()); + } + Session session = validator.getStored(ValidationStored.SessionTerminate); + rollbackExecutions(session); + rollbackAssetBalances(session.getId()); + deleteLiabilitiesAndClaims(session.getId()); + sessionMng.endSession(session); + return reqInfo.success(session.getId()); + } + + private void rollbackExecutions(Session session) { + Instant now = Instant.now(); + + Consumer> nullifier = imdg -> { + Collection execs = imdg.getCollectionObjectsBySQL( + "sessionId = " + session.getId() + ); + for (ExecutionCommon exec : execs) { + exec.setSessionId(null); + exec.setCoverageStatus(null); + exec.setUpdated(now); + imdg.update(exec); + } + log.info("updated {} executions with sessionId={}", execs.size(), session.getId()); + }; + log.debug("session.id={} section {}", session.getId(), session.getSection()); + + if (Section.FOND.equalsByKey(session.getSection())) { + nullifier.accept(upcastImdg(execFondImdg)); + } else if (Section.MKR.equalsByKey(session.getSection())) { + nullifier.accept(upcastImdg(execDpstImdg)); + } else throw new IllegalStateException("unknown section."); + } + + private void rollbackAssetBalances(Long sessionId) { + Instant now = Instant.now(); + ImdgPredicateBuilder pb = rgsImdg.predicateBuilder(); + ImdgPredicate prdct = pb.and( + pb.equals("sessionId", sessionId), + pb.sql(RegistryCodeSqlBuilder.getInstance(RegistryTradingParams.A__B).build()), + pb.not(pb.equals("balance", BigDecimal.ZERO)) + ); + + Collection a__bs = rgsImdg.getCollectionObjectsByPredicate(prdct); + a__bs.forEach(rgs -> { + rgs.setBalance(BigDecimal.ZERO); + rgs.setUpdated(now); + rgsImdg.update(rgs); + assets.processByAm_b(rgs, BigDecimal.ZERO); + }); + log.debug("set ZERO balance for {} A**B by predicate: {}", a__bs.size(), prdct); + } + + private void deleteLiabilitiesAndClaims(Long sessionId) { + ImdgPredicateBuilder pb = rgsImdg.predicateBuilder(); + ImdgPredicate prdct = pb.and( + pb.equals("sessionId", sessionId), + pb.or( + pb.sql(RegistryCodeSqlBuilder.getInstance(RegistryTradingParams.C__T).build()), + pb.sql(RegistryCodeSqlBuilder.getInstance(RegistryTradingParams.L__T).build()), + pb.sql(RegistryCodeSqlBuilder.getInstance(RegistryTradingParams.O__T).build()), + pb.sql(RegistryCodeSqlBuilder.getInstance(RegistryTradingParams.T__T).build()) + ) + ); + Collection liabilitiesAndClaims = rgsImdg.getCollectionObjectsByPredicate(prdct); + liabilitiesAndClaims.forEach(rgsImdg::delete); + log.debug("deleted {} liabilities and claims by predicate: {}", liabilitiesAndClaims.size(), prdct); + } + + + private static Imdg upcastImdg(Imdg imdg) { + @SuppressWarnings("unchecked") + Imdg result = (Imdg) imdg; + return result; + } +} 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 17f5b71e1..f3358101a 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 @@ -104,6 +104,9 @@ public record RegistryTradingParams(RegistryDesignation registryDesignation, public final static RegistryTradingParams AM__; public final static RegistryTradingParams AS__; public final static RegistryTradingParams CM__; + public final static RegistryTradingParams C__T; + public final static RegistryTradingParams O__T; + public final static RegistryTradingParams T__T; public final static RegistryTradingParams LM__; public final static RegistryTradingParams DMAU; public final static RegistryTradingParams DMAI; @@ -240,6 +243,18 @@ public record RegistryTradingParams(RegistryDesignation registryDesignation, null, null, RegistryUnit.V); + C__T = new RegistryTradingParams(RegistryDesignation.C, + null, + null, + RegistryUnit.T); + O__T = new RegistryTradingParams(RegistryDesignation.O, + null, + null, + RegistryUnit.T); + T__T = new RegistryTradingParams(RegistryDesignation.T, + null, + null, + RegistryUnit.T); } } diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java index c55d21d98..82bf2aaa1 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java @@ -139,6 +139,7 @@ public interface Consts { String SDF56_PROCESS = "sdf56-process"; String CONTINUE_SESSION_BN_FIRST_PART = "sdf57-process"; String CONTINUE_SESSION_BN_SECOND_PART = "sdf13-process"; + String KILL_SESSION = "kill-session"; String SWT_EXPORTER = "swt-exporter"; String REVISE_PROCESS = "revise-process"; String EXPORT_PROCESS = "export-process";