From 1bb865e32c2929fed91d752afc40c2c14e40b081 Mon Sep 17 00:00:00 2001 From: ialbert Date: Mon, 23 Sep 2024 18:59:33 +0300 Subject: [PATCH] =?UTF-8?q?=D0=BF=D1=80=D0=B8=D0=BC=D0=B5=D1=80=20spring?= =?UTF-8?q?=20state=20machine=20=D0=B4=D0=BB=D1=8F=20FinalSession,=20?= =?UTF-8?q?=D1=82=D0=B5=D1=81=D1=82,=20=D0=BF=D0=BE=D0=BD=D0=B8=D0=B7?= =?UTF-8?q?=D0=B8=D0=BB=20=D0=B2=D0=B5=D1=80=D1=81=D0=B8=D1=8E?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- clearing-parent/clearing-service/pom.xml | 2 +- .../MachineTestService.java | 21 +- .../SessionEvent.java | 2 +- .../StageDataEnum.java | 2 +- .../StateActionAdapter.java | 7 +- .../StateBnConfig.java | 27 +- .../FinalSessionStateMachineConfig.java | 173 +++++++ .../clearing/session/state/ActionHeader.java | 5 + .../clearing/session/state/DataEnum.java | 9 + .../clearing/session/state/SsnEvent.java | 7 + ...stractSessionActionForOkErrorHandling.java | 36 ++ .../state/action/CreateSessionAction.java | 43 ++ .../state/action/DealsPrepareAction.java | 140 ++++++ .../state/action/InclusionToPoolAction.java | 245 ++++++++++ ...pectionObligationsDepositReturnAction.java | 285 +++++++++++ .../action/InspectionObligationsV2Action.java | 408 ++++++++++++++++ .../action/ObligationAdmissionAction.java | 135 ++++++ ...equirementAndObligationCreationAction.java | 454 ++++++++++++++++++ .../session/state/action/ReviseStage1.java | 55 +++ .../session/state/action/Sdf56Action.java | 97 ++++ .../factory/SessionStateMachineWrapper.java | 110 +++++ .../SessionStatusChangingInterceptor.java | 108 +++++ .../state/listener/MachineStopListener.java | 43 ++ .../spcex/clearing/util/StateMachineUtil.java | 28 ++ .../session/teststate/StateMachineTest.java | 111 +++++ .../session/teststate/config/Event.java | 5 + .../session/teststate/config/State.java | 8 + .../config/TestStateMachineConfig.java | 150 ++++++ .../TestStateMachineExecutorsConfig.java | 32 ++ .../action/ContinueReviseTestAction.java | 21 + .../config/action/InitialAction.java | 21 + .../teststate/config/action/SelfAction.java | 21 + .../config/action/Send56TestAction.java | 16 + .../config/guard/NoActiveSessionGuard.java | 17 + .../src/test/resources/logback-test.xml | 19 + 35 files changed, 2840 insertions(+), 23 deletions(-) rename clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/{session => state_machine_1}/MachineTestService.java (55%) rename clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/{session => state_machine_1}/SessionEvent.java (54%) rename clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/{session => state_machine_1}/StageDataEnum.java (64%) rename clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/{session => state_machine_1}/StateActionAdapter.java (94%) rename clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/{session => state_machine_1}/StateBnConfig.java (94%) create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_2/specific/FinalSessionStateMachineConfig.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/ActionHeader.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/DataEnum.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/SsnEvent.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/AbstractSessionActionForOkErrorHandling.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/CreateSessionAction.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/DealsPrepareAction.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/InclusionToPoolAction.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/InspectionObligationsDepositReturnAction.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/InspectionObligationsV2Action.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/ObligationAdmissionAction.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/RequirementAndObligationCreationAction.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/ReviseStage1.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/Sdf56Action.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/factory/SessionStateMachineWrapper.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/interceptor/SessionStatusChangingInterceptor.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/listener/MachineStopListener.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/util/StateMachineUtil.java create mode 100644 clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/StateMachineTest.java create mode 100644 clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/Event.java create mode 100644 clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/State.java create mode 100644 clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/TestStateMachineConfig.java create mode 100644 clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/TestStateMachineExecutorsConfig.java create mode 100644 clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/action/ContinueReviseTestAction.java create mode 100644 clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/action/InitialAction.java create mode 100644 clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/action/SelfAction.java create mode 100644 clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/action/Send56TestAction.java create mode 100644 clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/guard/NoActiveSessionGuard.java create mode 100644 clearing-parent/clearing-service/src/test/resources/logback-test.xml diff --git a/clearing-parent/clearing-service/pom.xml b/clearing-parent/clearing-service/pom.xml index 77c99ea1e..3fa457dd6 100644 --- a/clearing-parent/clearing-service/pom.xml +++ b/clearing-parent/clearing-service/pom.xml @@ -20,7 +20,7 @@ org.springframework.statemachine spring-statemachine-starter - 3.2.0 + 2.5.1 ru.spcex.clearing diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/session/MachineTestService.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_1/MachineTestService.java similarity index 55% rename from clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/session/MachineTestService.java rename to clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_1/MachineTestService.java index fe980bed3..8c727246d 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/session/MachineTestService.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_1/MachineTestService.java @@ -1,9 +1,11 @@ -package ru.spcex.clearing.config.session; +package ru.spcex.clearing.config.state_machine_1; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.InitializingBean; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.messaging.Message; +import org.springframework.messaging.support.MessageBuilder; import org.springframework.statemachine.StateMachine; import org.springframework.statemachine.config.StateMachineFactory; import org.springframework.stereotype.Component; @@ -23,11 +25,16 @@ public class MachineTestService implements InitializingBean { @Override public void afterPropertiesSet() throws Exception { -// stateMachine.start(); -// log.info("send start revise event"); -// stateMachine.sendEvent(SessionEvent.Revise); -// log.info("send continue revise event"); -// stateMachine.sendEvent(SessionEvent.SdfReceived); -// log.info("end test"); + stateMachine.start(); + log.info("send start revise event"); + stateMachine.sendEvent(SessionEvent.Revise); + log.info("send continue revise event"); + stateMachine.sendEvent(SessionEvent.SdfReceived); + Message eventMessage = MessageBuilder + .withPayload(SessionEvent.SdfReceived) + .setHeader("myPayload", new Object()) // Set the payload in the message header + .build(); + stateMachine.sendEvent(eventMessage); + log.info("end test"); } } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/session/SessionEvent.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_1/SessionEvent.java similarity index 54% rename from clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/session/SessionEvent.java rename to clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_1/SessionEvent.java index 50430df1a..7006d6bee 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/session/SessionEvent.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_1/SessionEvent.java @@ -1,4 +1,4 @@ -package ru.spcex.clearing.config.session; +package ru.spcex.clearing.config.state_machine_1; public enum SessionEvent { Revise, diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/session/StageDataEnum.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_1/StageDataEnum.java similarity index 64% rename from clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/session/StageDataEnum.java rename to clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_1/StageDataEnum.java index c12287a21..e47941657 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/session/StageDataEnum.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_1/StageDataEnum.java @@ -1,4 +1,4 @@ -package ru.spcex.clearing.config.session; +package ru.spcex.clearing.config.state_machine_1; public enum StageDataEnum { sessionId, session, ExecutionList, PaymentInstructions diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/session/StateActionAdapter.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_1/StateActionAdapter.java similarity index 94% rename from clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/session/StateActionAdapter.java rename to clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_1/StateActionAdapter.java index e0099ecfd..e04250186 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/session/StateActionAdapter.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_1/StateActionAdapter.java @@ -1,5 +1,6 @@ -package ru.spcex.clearing.config.session; +package ru.spcex.clearing.config.state_machine_1; +import java.util.function.Function; import org.springframework.statemachine.ExtendedState; import org.springframework.statemachine.StateContext; import org.springframework.statemachine.action.Action; @@ -8,8 +9,6 @@ import ru.spcex.clearing.session.stage.StageResult; import ru.spcex.clearing.session.stage.Task; import ru.spcex.clearing.session.stage.TaskType; -import java.util.function.Function; - /** * адаптеры между ISessionStage старого образца к Spring State Machine
* главное назначение: это сохранение результатов работы старого ISessionStage в ExtendedState
@@ -20,7 +19,7 @@ public class StateActionAdapter implements Action { private Function payloadForStageGetter; private StageDataEnum saveName; - protected StateActionAdapter(ISessionStage session) { + public StateActionAdapter(ISessionStage session) { this.session = session; } 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/state_machine_1/StateBnConfig.java similarity index 94% rename from clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/session/StateBnConfig.java rename to clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_1/StateBnConfig.java index 7f685f721..bc7cefd3c 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/state_machine_1/StateBnConfig.java @@ -1,5 +1,11 @@ -package ru.spcex.clearing.config.session; +package ru.spcex.clearing.config.state_machine_1; +import java.time.LocalDate; +import java.util.Arrays; +import java.util.HashSet; +import java.util.List; +import java.util.Map; +import java.util.function.Supplier; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Qualifier; @@ -15,7 +21,17 @@ 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.session.stage.TaskType; -import ru.spcex.clearing.session.stage.impl.*; +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.FormingPaymentInstruction; +import ru.spcex.clearing.session.stage.impl.FormingRegistersOnOS; +import ru.spcex.clearing.session.stage.impl.InclusionObligations; +import ru.spcex.clearing.session.stage.impl.InspectionObligations; +import ru.spcex.clearing.session.stage.impl.ObligationAdmission; +import ru.spcex.clearing.session.stage.impl.RequirementsAndObligationCreation; +import ru.spcex.clearing.session.stage.impl.UnlockResources; import ru.spcex.clearing.session.stage.task.DealsPreparePayload; import ru.spcex.clearing.session.stage.task.FormingPaymentInstructionPayload; import ru.spcex.clearing.session.stage.task.InclusionToPoolPayload; @@ -28,13 +44,6 @@ import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.ImdgProvider; import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder; -import java.time.LocalDate; -import java.util.Arrays; -import java.util.HashSet; -import java.util.List; -import java.util.Map; -import java.util.function.Supplier; - @Configuration("sessionStateMachineFactory") @EnableStateMachineFactory(name = "PrimaryBnSessionStateMachineFactory") public class StateBnConfig extends EnumStateMachineConfigurerAdapter { diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_2/specific/FinalSessionStateMachineConfig.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_2/specific/FinalSessionStateMachineConfig.java new file mode 100644 index 000000000..84b40c35d --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_2/specific/FinalSessionStateMachineConfig.java @@ -0,0 +1,173 @@ +package ru.spcex.clearing.config.state_machine_2.specific; + +import java.time.LocalDate; +import java.util.Arrays; +import java.util.HashSet; +import java.util.List; +import java.util.function.Supplier; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.context.annotation.Configuration; +import org.springframework.statemachine.config.EnableStateMachineFactory; +import org.springframework.statemachine.config.EnumStateMachineConfigurerAdapter; +import org.springframework.statemachine.config.builders.StateMachineConfigurationConfigurer; +import org.springframework.statemachine.config.builders.StateMachineStateConfigurer; +import org.springframework.statemachine.config.builders.StateMachineTransitionConfigurer; +import org.springframework.statemachine.listener.StateMachineListenerAdapter; +import org.springframework.statemachine.state.State; +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.session.stage.TaskType; +import ru.spcex.clearing.session.state.SsnEvent; +import ru.spcex.clearing.session.state.action.CreateSessionAction; +import ru.spcex.clearing.session.state.action.DealsPrepareAction; +import ru.spcex.clearing.session.state.action.InclusionToPoolAction; +import ru.spcex.clearing.session.state.action.InspectionObligationsDepositReturnAction; +import ru.spcex.clearing.session.state.action.InspectionObligationsV2Action; +import ru.spcex.clearing.session.state.action.ObligationAdmissionAction; +import ru.spcex.clearing.session.state.action.RequirementAndObligationCreationAction; +import ru.spcex.clearing.session.state.action.ReviseStage1; +import ru.spcex.clearing.session.state.action.Sdf56Action; +import ru.spcex.clearing.util.StateMachineUtil; +import ru.spcex.platform.classes.base.interfaces.ExecutionType; +import ru.spcex.platform.enumeration.RegistryStatus; +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; + +@Configuration("finalSessionStateMachineFactoryConfig") +@EnableStateMachineFactory(name = "FINL") +public class FinalSessionStateMachineConfig + extends EnumStateMachineConfigurerAdapter { + private final Logger log = LoggerFactory.getLogger(getClass()); + + private final Imdg sessionImdg; + private final Imdg rgsImdg; + private final Sdf56Action sendSdf56Action; + private final ReviseStage1 reviseStage1Action; + private final DealsPrepareAction dealsPrepareAction; + private final RequirementAndObligationCreationAction reqAndOblAction; + private final ObligationAdmissionAction obligationAdmissionAction; + private final InclusionToPoolAction inclusionToPoolAction; + private final InspectionObligationsDepositReturnAction inspOblDepositReturnAction; + private final InspectionObligationsV2Action inspectionObligationsV2Action; + + + @Autowired + public FinalSessionStateMachineConfig( + ImdgProvider imdgProvider, + @Qualifier("marketCodesForBn") + Supplier> marketCodes, + Sdf56Action sendSdf56Action, ReviseStage1 reviseStage1Action, DealsPrepareAction dealsPrepareAction, RequirementAndObligationCreationAction reqAndOblAction, ObligationAdmissionAction obligationAdmissionAction, InclusionToPoolAction inclusionToPoolAction, InspectionObligationsDepositReturnAction inspOblDepositReturnAction, InspectionObligationsV2Action inspectionObligationsV2Action + ) { + this.sendSdf56Action = sendSdf56Action; + this.sessionImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Session, Session.class); + this.rgsImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class); + this.reviseStage1Action = reviseStage1Action; + this.dealsPrepareAction = dealsPrepareAction; + this.reqAndOblAction = reqAndOblAction; + this.obligationAdmissionAction = obligationAdmissionAction; + this.inclusionToPoolAction = inclusionToPoolAction; + this.inspOblDepositReturnAction = inspOblDepositReturnAction; + this.inspectionObligationsV2Action = inspectionObligationsV2Action; + + //stages settings: + ImdgPredicateBuilder rgsPrctBuilder = rgsImdg.predicateBuilder(); + this.dealsPrepareAction.searchForExecutions(ExecutionType.ExecutionDeposit); + this.inclusionToPoolAction.addRegistryCondition( + rgsPrctBuilder.and( + rgsPrctBuilder.or( + rgsPrctBuilder.equals("registryStatus", RegistryStatus.PROC.getKey()), + rgsPrctBuilder.equals("registryStatus", RegistryStatus.MNG.getKey()) + ), + rgsPrctBuilder.lessEqual("valueDate", LocalDate.now()) + ) + ); + this.inspOblDepositReturnAction.setSessionType(SessionType.FINL); + this.inspectionObligationsV2Action.setSessionType(SessionType.FINL); + } + + @Override + public void configure(StateMachineStateConfigurer states) throws Exception { + states + .withStates() + .initial(TaskType.StartRevise, StateMachineUtil.chain(new CreateSessionAction(sessionImdg), sendSdf56Action)) + //.state(TaskType.StartRevise, sendSdf56Action) //entry action + .states(new HashSet<>(Arrays.asList( + TaskType.StartRevise, + TaskType.StartRevisePart1, + TaskType.DealsPrepare, + TaskType.RequirementsAndObligationsCreate, + TaskType.ObligationsAdmission, + TaskType.InclusionToPool, + TaskType.InspectionObligations, + TaskType.InspectionObligations, + TaskType.FormingRegistersOnOS, + TaskType.FormingPaymentInstruction, + TaskType.FormingPaymentInstruction, + TaskType.FinishingSession, + TaskType.EndStageNotification + ))) + .end(TaskType.EndStageNotification); + } + + + @Override + public void configure(StateMachineTransitionConfigurer transitions) throws Exception { + transitions + .withExternal() + .event(SsnEvent.Sdf57Processed) + .source(TaskType.StartRevise).target(TaskType.StartRevisePart1) + .action(reviseStage1Action) + .and() + .withExternal() + .source(TaskType.StartRevisePart1).target(TaskType.DealsPrepare) + .action(dealsPrepareAction) + .and() + .withExternal() + .source(TaskType.DealsPrepare).target(TaskType.RequirementsAndObligationsCreate) + .action(reqAndOblAction) + .and() + .withExternal() + .source(TaskType.RequirementsAndObligationsCreate).target(TaskType.ObligationsAdmission) + .action(obligationAdmissionAction) + .and() + .withExternal() + .source(TaskType.ObligationsAdmission).target(TaskType.InclusionToPool) + .action(inclusionToPoolAction) + .and() + .withExternal() + .source(TaskType.InclusionToPool).target(TaskType.InspectionObligations) + .action(inspOblDepositReturnAction) + .and() + //internal (no state change) triggerless(without an event) action + .withInternal() + .source(TaskType.InspectionObligations) + .action(inspectionObligationsV2Action); + + + + //TaskType.RequirementsAndObligationsCreate + + } + + @Override + public void configure(StateMachineConfigurationConfigurer config) throws Exception { + StateMachineListenerAdapter loggingChangeStateListener = new StateMachineListenerAdapter<>() { + @Override + public void stateEntered(State state) { + TaskType enteredState = state != null ? state.getId() : null; + log.info("State entered: {}", enteredState); + } + }; + config + .withConfiguration() + .machineId(SessionType.FINL.getKey()) + .listener(loggingChangeStateListener) + ; + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/ActionHeader.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/ActionHeader.java new file mode 100644 index 000000000..78b1c0724 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/ActionHeader.java @@ -0,0 +1,5 @@ +package ru.spcex.clearing.session.state; + +public enum ActionHeader { + sdf56FromTime; +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/DataEnum.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/DataEnum.java new file mode 100644 index 000000000..42c3946b6 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/DataEnum.java @@ -0,0 +1,9 @@ +package ru.spcex.clearing.session.state; + +public enum DataEnum { + sessionId, //Long + sessionType, //Session + session, //Session + dealsPrepared, //List + counterPartyId, //Long +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/SsnEvent.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/SsnEvent.java new file mode 100644 index 000000000..ff471b212 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/SsnEvent.java @@ -0,0 +1,7 @@ +package ru.spcex.clearing.session.state; + +public enum SsnEvent { + Revise, + Sdf57Processed, + SdfReceived, +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/AbstractSessionActionForOkErrorHandling.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/AbstractSessionActionForOkErrorHandling.java new file mode 100644 index 000000000..aa8459a3c --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/AbstractSessionActionForOkErrorHandling.java @@ -0,0 +1,36 @@ +package ru.spcex.clearing.session.state.action; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.statemachine.StateContext; +import org.springframework.statemachine.action.Action; +import ru.clearing.classes.statics.data.misc.Session; +import ru.spcex.clearing.session.stage.TaskType; +import ru.spcex.clearing.session.state.DataEnum; +import ru.spcex.clearing.session.state.SsnEvent; + +public abstract class AbstractSessionActionForOkErrorHandling implements Action { + private final Logger log = LoggerFactory.getLogger(getClass()); + + @Override + public void execute(StateContext context) { + try { + Long sessionId = context.getExtendedState().get(DataEnum.sessionId, Long.class); + Session session = context.getExtendedState().get(DataEnum.session, Session.class); + if (sessionId != null) { + log.info("sessionId={}{}: {} action started", + sessionId, + session != null ? ("/%s".formatted(session.getSessionType())) : "", + getClass().getSimpleName()); + } else { + log.info("{} action started", getClass().getSimpleName()); + } + actualExecute(context); + } catch (Exception e) { + context.getStateMachine().setStateMachineError(e); + throw e; + } + } + + protected abstract void actualExecute(StateContext context); +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/CreateSessionAction.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/CreateSessionAction.java new file mode 100644 index 000000000..944537ecd --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/CreateSessionAction.java @@ -0,0 +1,43 @@ +package ru.spcex.clearing.session.state.action; + +import java.time.LocalDate; +import java.util.Map; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.statemachine.StateContext; +import ru.clearing.classes.statics.data.misc.Session; +import ru.spcex.clearing.session.stage.TaskType; +import ru.spcex.clearing.session.state.DataEnum; +import ru.spcex.clearing.session.state.SsnEvent; +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; + +public class CreateSessionAction extends AbstractSessionActionForOkErrorHandling { + private final Imdg sessionImdg; + private final Logger log = LoggerFactory.getLogger(CreateSessionAction.class); + + public CreateSessionAction(Imdg sessionImdg) { + this.sessionImdg = sessionImdg; + } + + @Override + public void actualExecute(StateContext context) { + Session existActiveSession = sessionImdg.getFirstObjectByFieldValues(Map.of("workflowStatus", SessionStatus.ACTV.getKey())); + if (existActiveSession != null) { + throw new RuntimeException("ActiveSessionIsPresent " + existActiveSession); + } + Session newSession = new Session(); + newSession.setSection(Section.FOND.getKey());//fixme + newSession.setSessionType(SessionType.IPOB.getKey());//fixme + newSession.setSessionStatus(TaskType.StartRevise.getKey()); + newSession.setWorkflowStatus(SessionStatus.ACTV.getKey()); + newSession.setClearingDate(LocalDate.now()); + sessionImdg.insert(newSession); + log.info("started new session.id={}", newSession.getId()); + context.getExtendedState().getVariables().put(DataEnum.sessionId.name(), newSession.getId()); + context.getExtendedState().getVariables().put(DataEnum.sessionType.name(), SessionType.IPOB); + context.getExtendedState().getVariables().put(DataEnum.session.name(), newSession); + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/DealsPrepareAction.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/DealsPrepareAction.java new file mode 100644 index 000000000..ad6767b9e --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/DealsPrepareAction.java @@ -0,0 +1,140 @@ +package ru.spcex.clearing.session.state.action; + +import java.time.LocalDate; +import java.util.ArrayList; +import java.util.List; +import java.util.function.BiFunction; +import java.util.function.Function; +import java.util.stream.Collectors; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.config.ConfigurableBeanFactory; +import org.springframework.context.annotation.Scope; +import org.springframework.statemachine.StateContext; +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.ExecutionDeposit; +import ru.clearing.classes.statics.data.execution.ExecutionFond; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.session.stage.TaskType; +import ru.spcex.clearing.session.state.DataEnum; +import ru.spcex.clearing.session.state.SsnEvent; +import ru.spcex.platform.classes.base.interfaces.ExecutionType; +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; + +@Service +@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) +public class DealsPrepareAction extends AbstractSessionActionForOkErrorHandling { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final Imdg executionDepositImdg; + private final Imdg executionFondImdg; + private final Imdg executionCurrencyImdg; + private final List execDepositPredicates = new ArrayList<>(); + private final List execFondPredicates = new ArrayList<>(); + private final List execCurrencyPredicates = new ArrayList<>(); + private ExecutionType executionType = ExecutionType.ExecutionFond; + + @Autowired + public DealsPrepareAction(ImdgProvider imdgProvider) { + this.executionDepositImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionDeposit, ExecutionDeposit.class); + this.executionFondImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionFond, ExecutionFond.class); + this.executionCurrencyImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionCurrency, ExecutionCurrency.class); + } + + public void searchForExecutions(ExecutionType executionType) { + this.executionType = executionType; + } + + @Override + public void actualExecute(StateContext ctx) { + Long sessionId = ctx.getExtendedState().get(DataEnum.sessionId, Long.class); + Long companyId = null, securityId = null; //никогда не использовалось + + log.info("loading deals sessionId: {}, companyId: {}, securityId: {}", sessionId, companyId, securityId); + StringBuilder executionDepositSQL = new StringBuilder("sessionId = null"); + if (companyId != null) { + executionDepositSQL.append(" and companyId = ").append(companyId); //fixme вероятно в ТЗ недостаточно условий для выборки - counterPartyId + } + if (securityId != null) { + executionDepositSQL.append(" and securityId = ").append(securityId); + } + ImdgPredicateBuilder prdctBuilder = executionFondImdg.predicateBuilder(); + + ImdgPredicate sqlPredicate = prdctBuilder.and( + prdctBuilder.sql(executionDepositSQL.toString()), + prdctBuilder.equals("tradingDate", LocalDate.now()) + + ); + BiFunction, Imdg, ImdgPredicate> prdComposer = (imdgPredicates, imdg) -> { + ImdgPredicateBuilder pb = imdg.predicateBuilder(); + if (imdgPredicates.isEmpty()) { + return sqlPredicate; + } else { + return pb.and(sqlPredicate, pb.and(imdgPredicates.toArray(new ImdgPredicate[0]))); + } + }; + List excs; + switch (executionType) { + case ExecutionDeposit -> { + ImdgPredicate excDepPrct = prdComposer.apply(execDepositPredicates, executionDepositImdg); + log.info("using predicate to load deals: {}", excDepPrct.toString()); + excs = executionDepositImdg.getCollectionObjectsByPredicate(excDepPrct) + .stream() + .map(execToInterface()) + .collect(Collectors.toList()); + } + case ExecutionFond -> { + ImdgPredicate excFondPrct = prdComposer.apply(execFondPredicates, executionFondImdg); + log.info("using predicate to load deals: {}", excFondPrct.toString()); + excs = executionFondImdg.getCollectionObjectsByPredicate(excFondPrct) + .stream() + .map(execToInterface()) + .collect(Collectors.toList()); + } + case ExecutionCurrency -> { + ImdgPredicate excCurrPrct = prdComposer.apply(execCurrencyPredicates, executionCurrencyImdg); + log.info("using predicate to load deals: {}", excCurrPrct.toString()); + excs = executionCurrencyImdg.getCollectionObjectsByPredicate(excCurrPrct) + .stream() + .map(execToInterface()) + .collect(Collectors.toList()); + } + default -> throw new IllegalStateException("Unknown execution type: " + executionType); + } + for (ExecutionCommon exc : excs) { + exc.setSessionId(sessionId); + if (exc instanceof ExecutionDeposit) { + executionDepositImdg.update((ExecutionDeposit) exc); + } else if (exc instanceof ExecutionFond) { + executionFondImdg.update((ExecutionFond) exc); + } else if (exc instanceof ExecutionCurrency) { + executionCurrencyImdg.update((ExecutionCurrency) exc); + } + } + ctx.getExtendedState().getVariables().put(DataEnum.dealsPrepared, excs); + } + + public DealsPrepareAction addExecutionDepositCondition(ImdgPredicate imdgPredicate) { + this.execDepositPredicates.add(imdgPredicate); + return this; + } + + public DealsPrepareAction addExecutionFondCondition(ImdgPredicate imdgPredicate) { + this.execFondPredicates.add(imdgPredicate); + return this; + } + + public DealsPrepareAction addExecutionCurrencyCondition(ImdgPredicate imdgPredicate) { + this.execCurrencyPredicates.add(imdgPredicate); + return this; + } + + private static Function execToInterface() { + return (e) -> e; + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/InclusionToPoolAction.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/InclusionToPoolAction.java new file mode 100644 index 000000000..f54b2d63e --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/InclusionToPoolAction.java @@ -0,0 +1,245 @@ +package ru.spcex.clearing.session.state.action; + +import java.time.Instant; +import java.time.LocalDate; +import java.util.ArrayList; +import java.util.Collection; +import java.util.Collections; +import java.util.Comparator; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.Objects; +import java.util.Optional; +import java.util.stream.Collectors; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.config.ConfigurableBeanFactory; +import org.springframework.context.annotation.Scope; +import org.springframework.statemachine.StateContext; +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.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.platform.messaging.service.sender.KafkaSender; +import ru.spcex.clearing.session.stage.TaskType; +import ru.spcex.clearing.session.state.DataEnum; +import ru.spcex.clearing.session.state.SsnEvent; +import ru.spcex.clearing.util.ClearingUtil; +import ru.spcex.platform.classes.base.interfaces.ExecutionType; +import ru.spcex.platform.enumeration.RegistryDesignation; +import ru.spcex.platform.enumeration.RegistryInstrumentType; +import ru.spcex.platform.enumeration.RegistryStatus; +import ru.spcex.platform.enumeration.RegistryUnit; +import ru.spcex.platform.enumeration.Section; +import ru.spcex.platform.enumeration.SessionType; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgId; +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.utils.enumeration.IEnumKey; +import ru.spcex.platform.utils.enumeration.IMessageResolver; +import ru.spcex.platform.utils.text.TextUtil; + +@Service +@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) +public class InclusionToPoolAction extends AbstractSessionActionForOkErrorHandling { + private final Logger log = LoggerFactory.getLogger(getClass()); + //todo remove (set all in single method setImdg(provider -> setImdg1();setIdGenerator();...) + private ImdgProvider imdgProvider; + private ImdgId idGenerator; + private Imdg registryImdg; + private Imdg sessionImdg; + private final Imdg executionDepositImdg; + private final Imdg executionFondImdg; + private final Imdg executionCurrImdg; + private KafkaSender kafkaSender; + private final List registryConditions = new ArrayList<>(); + private SessionType sessionType; + private final ImdgPredicateBuilder rgsPb; + + @Autowired + public InclusionToPoolAction(ImdgProvider imdgProvider, IMessageResolver messageResolver) { + this.imdgProvider = imdgProvider; + this.idGenerator = imdgProvider.getImdgIdGenerator(); + this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class); + this.sessionImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Session, Session.class); + this.executionDepositImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionDeposit, ExecutionDeposit.class); + this.executionFondImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionFond, ExecutionFond.class); + this.executionCurrImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionCurrency, ExecutionCurrency.class); + this.rgsPb = registryImdg.predicateBuilder(); + } + + /** + * для сессии возврата депозитов статус регистров + * proc ИЛИ mng + * proc не захардкожен для случая дополнительных условий + * указывать обязательно если есть дополнительные условия + */ + public InclusionToPoolAction addRegistryCondition(ImdgPredicate imdgPredicate) { + this.registryConditions.add(imdgPredicate); + return this; + } + + @Override + protected void actualExecute(StateContext ctx) { + Long sessionId = ctx.getExtendedState().get(DataEnum.sessionId, Long.class); + Long counterPartyId = ctx.getExtendedState().get(DataEnum.counterPartyId, Long.class); + SessionType sessionType = ctx.getExtendedState().get(DataEnum.sessionType, SessionType.class); + String sqlCondition = String.format("registryDesignation in ('%s', '%s') and " + + "registryInstrumentType in ('%s', '%s') and " + + "registryUnit = '%s'", + RegistryDesignation.O.getKey(), RegistryDesignation.T.getKey(), + RegistryInstrumentType.S.getKey(), RegistryInstrumentType.M.getKey(), + RegistryUnit.T.getKey()); + ImdgPredicateBuilder prdctBuilder = registryImdg.predicateBuilder(); + ImdgPredicate defaultCondition = prdctBuilder.sql(sqlCondition); + if (sessionType != null) { + //не всегда нужное условие + defaultCondition = prdctBuilder.and(defaultCondition, sessionTypePredicate()); + } + ImdgPredicate registryPredicate; + if (registryConditions.size() > 0) { + registryPredicate = prdctBuilder.and( + defaultCondition, + prdctBuilder.and(registryConditions.toArray(new ImdgPredicate[0])) + ); + } else { + registryPredicate = prdctBuilder.and( + defaultCondition, + SessionType.UNIT.equals(sessionType) ? prdctBuilder.equals("sessionId", sessionId) : prdctBuilder.alwaysTrue(), + prdctBuilder.equals("registryStatus", RegistryStatus.PROC.getKey()) + ); + } + log.info("Inclusion to pool with predicate: {}", registryPredicate.toString()); + Collection obligations = registryImdg.getCollectionObjectsByPredicate(registryPredicate); + Map> registryByGroupId = obligations.stream() + .filter(registry -> registry.getSettlementDate().isEqual(LocalDate.now())) + .collect(Collectors.groupingBy(Registry::getGroupId)); + + Optional ssnType = Optional.of(sessionType); + Map rgsToUpdate = new HashMap<>(); + List execsToUpdate = new ArrayList<>(); + boolean loadExecs = !IEnumKey.contains(sessionType, + SessionType.TRDT, SessionType.CURR, SessionType.UNIT, SessionType.IPOT); + for (Map.Entry> entrySet : registryByGroupId.entrySet()) { + log.debug("Processing set of registry with groupId: {}", entrySet.getKey()); + String rgsSection = null; + for (Registry registry : entrySet.getValue()) { + registry.setRegistryStatus(RegistryStatus.POOL.getKey()); + registry.setClearingDate(LocalDate.now()); + registry.setSessionId(sessionId); + ssnType.ifPresent(st -> registry.setSessionType(st.getKey())); + rgsToUpdate.put(registry.getId(), registry); + if (TextUtil.isEmpty(rgsSection)) { + rgsSection = registry.getSection(); + } + } + if (loadExecs) { + execsToUpdate.addAll(updateExecutions(sessionId, sessionType, entrySet.getKey(), rgsSection)); + } + } + registryImdg.putAll(rgsToUpdate); + if (loadExecs) { + updateExecsBatchV2(execsToUpdate); + } + } + + @SuppressWarnings("unchecked") + private Collection updateExecutions(Long sessionId, SessionType ssnTpe, Long rgsGroupId, String rgsSection) { + Instant now = Instant.now(); + if (ssnTpe == null || sessionId == null) { + log.debug("will not update executions#sessionId - couldn't determine session type"); + return Collections.emptyList(); + } + Imdg execImdg = null; + switch (ssnTpe) { + case MEDM, FINL, XDEP -> execImdg = (Imdg) executionDepositImdg; + case IPO0, IPOB, TRDT, IPOT -> execImdg = (Imdg) executionFondImdg; + case CURR -> execImdg = (Imdg) executionCurrImdg; + default -> { + //кейс для "общих" сессий (напр. UNIT) в которых сочетаются разные сделки + //смотрим на секцию регистра. + Section section; + if ((section = IEnumKey.getEnumByKey(Section.class, rgsSection)) != null) { + if (section.equals(Section.MKR)) execImdg = (Imdg) executionDepositImdg; + else if (section.equals(Section.FOND)) execImdg = (Imdg) executionFondImdg; + else if (section.equals(Section.CURR)) execImdg = (Imdg) executionCurrImdg; + } + } + } + if (execImdg == null) { + log.warn("couldn't define Execution Type for session {}. Will not update executions#sessionId", + this.sessionType); + return Collections.emptyList(); + } + ImdgPredicateBuilder pb = execImdg.predicateBuilder(); + ImdgPredicate prdct = pb.equals("exchangeExecutionId", rgsGroupId); + Collection execs = execImdg.getCollectionObjectsByPredicate(prdct); + //Imdg finalExecImdg = execImdg; + execs.forEach(e -> { + e.setUpdated(now); + e.setSessionId(sessionId); +// finalExecImdg.update(e); + }); + return (Collection) execs; + } + + private ImdgPredicate sessionTypePredicate() { + if (this.sessionType != null) { + if (SessionType.FINL.equals(sessionType) || SessionType.MEDM.equals(sessionType)) { + return rgsPb.or( + rgsPb.equals("sessionType", SessionType.FINL.getKey()), + rgsPb.equals("sessionType", SessionType.MEDM.getKey()), + rgsPb.equals("sessionType", SessionType.XDEP.getKey()) + ); + } else { + return rgsPb.equals("sessionType", sessionType.getKey()); + } + } else { + return rgsPb.alwaysTrue(); + } + } + + private void updateExecsBatchV2(List execs) { + log.debug("updating {} executions", execs.size()); + execs.sort(Comparator.comparing(ExecutionCommon::type)); + log.debug("sorted executions by type"); + + Map m = new HashMap<>(); + for (int i = 0; i < execs.size(); i++) { + ExecutionCommon exec = execs.get(i); + ExecutionType execType = exec.type(); + log.debug("batch from {} position type {}", i, execType); + m.put(exec.getId(), exec); + int j = i + 1; + while (j < execs.size() && j < i + 100) { + ExecutionCommon other = execs.get(j); + if (!Objects.equals(other.type(), execType)) { + break; + } + m.put(other.getId(), other); + j++; + } + i = j; + log.debug("batch size {} type {}", m.size(), execType); + if (!m.isEmpty()) { + switch (execType) { + case ExecutionDeposit -> executionDepositImdg.putAll(ClearingUtil.castMap(m)); + case ExecutionFond -> executionFondImdg.putAll(ClearingUtil.castMap(m)); + case ExecutionCurrency -> executionCurrImdg.putAll(ClearingUtil.castMap(m)); + } + } + log.debug("batch size {} type {} done", m.size(), execType); + m.clear(); + } + } + + +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/InspectionObligationsDepositReturnAction.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/InspectionObligationsDepositReturnAction.java new file mode 100644 index 000000000..51d3b2b61 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/InspectionObligationsDepositReturnAction.java @@ -0,0 +1,285 @@ +package ru.spcex.clearing.session.state.action; + +import java.math.BigDecimal; +import java.time.Instant; +import java.util.List; +import java.util.Map; +import java.util.Optional; +import java.util.function.BiConsumer; +import java.util.stream.Collectors; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.config.ConfigurableBeanFactory; +import org.springframework.context.annotation.Scope; +import org.springframework.statemachine.StateContext; +import org.springframework.stereotype.Service; +import ru.clearing.classes.statics.data.company.ClearingMemberCategory; +import ru.clearing.classes.statics.data.registry.Registry; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.service.AssetTrio; +import ru.spcex.clearing.service.integration.GatewayRequestCreator; +import ru.spcex.clearing.service.registry.AssetTBFProcessing; +import ru.spcex.clearing.service.registry.RegistryManager; +import ru.spcex.clearing.service.schedule.TradingTimeService; +import ru.spcex.clearing.session.stage.TaskType; +import ru.spcex.clearing.session.stage.impl.GatewayRequester; +import ru.spcex.clearing.session.stage.impl.PlanBalanceCalc; +import ru.spcex.clearing.session.state.DataEnum; +import ru.spcex.clearing.session.state.SsnEvent; +import ru.spcex.platform.enumeration.ClearingCategory; +import ru.spcex.platform.enumeration.RegistryCapacity; +import ru.spcex.platform.enumeration.RegistryDesignation; +import ru.spcex.platform.enumeration.RegistryInstrumentType; +import ru.spcex.platform.enumeration.RegistryStatus; +import ru.spcex.platform.enumeration.RegistryTradingParams; +import static ru.spcex.platform.enumeration.RegistryTradingParams.*; +import ru.spcex.platform.enumeration.RegistryUnit; +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.specific.RegistryCodeSqlBuilder; +import ru.spcex.platform.utils.enumeration.IEnumKey; +import static ru.spcex.platform.utils.number.BigDecimalUtil.safeBD; + +@Service +@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) +public class InspectionObligationsDepositReturnAction extends AbstractSessionActionForOkErrorHandling { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final Imdg registryImdg; + private final Imdg categoryImdg; + private SessionType sessionType; + private final RegistryManager rgsMng; + private final AssetTBFProcessing assets; + private final TradingTimeService tradingTimeService; + private final GatewayRequester gateway; + private final PlanBalanceCalc planBalanceCalc; + + @Autowired + public InspectionObligationsDepositReturnAction(ImdgProvider imdgProvider, RegistryManager rgsMng, AssetTBFProcessing assets, TradingTimeService tradingTimeService, GatewayRequester gateway, PlanBalanceCalc planBalanceCalc) { + this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class); + this.categoryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ClearingMemberCategory, ClearingMemberCategory.class); + this.rgsMng = rgsMng; + this.assets = assets; + this.tradingTimeService = tradingTimeService; + this.gateway = gateway; + this.planBalanceCalc = planBalanceCalc; + this.gateway.setName("InspectionObligationsDepositReturnAction|OM*T"); + } + + @Override + protected void actualExecute(StateContext ctx) { + Long sessionId = ctx.getExtendedState().get(DataEnum.sessionId, Long.class); + String sqlCondition = String.format("(%s) and registryStatus = '%s'", + RegistryCodeSqlBuilder.getInstance(OS_T, OM_T, TS_T, TM_T).build(), + RegistryStatus.POOL.getKey()); + + List>> registriesByGroupSorted = registryImdg.getCollectionObjectsBySQL(sqlCondition) + .stream() + .collect(Collectors.groupingBy(Registry::getGroupId)) + .entrySet() + .stream() + .sorted(Map.Entry.comparingByKey()) + .toList(); + log.info("registry groups found {}", registriesByGroupSorted.size()); + + for (Map.Entry> entry : registriesByGroupSorted) { + List group = entry.getValue(); + Optional omtInGroupO = group.stream().filter(registry -> equalByRgs(OM_T, registry)).findFirst(); + Optional tmtInGroupO = group.stream().filter(registry -> equalByRgs(TM_T, registry)).findFirst(); + if (omtInGroupO.isEmpty() || tmtInGroupO.isEmpty()) { + log.error("groupId {} failed to find OM*T/TM*T registry", entry.getKey()); + group.forEach(rgs -> updateStatus(rgs, registryStatusFailed())); // FAIL or MNG + continue; + } + Registry omtRgs = omtInGroupO.get(); + Registry tmtRgs = tmtInGroupO.get(); + log.debug("groupId={}, OM*T.id={}", entry.getKey(), omtRgs.getId()); + BigDecimal omtBalance = safeBD(omtRgs.getBalance()); + + //первая часть сделки депозита + if (omtRgs.getValueDate() == null + || omtRgs.getSettlementDate() == null + || !omtRgs.getValueDate().isBefore(omtRgs.getSettlementDate())) { + continue; + } + + //вторая часть сделки - возврат депозита + Optional dmx = rgsMng.searchDmxByCounterParty(omtRgs); + if (dmx.isPresent()) { + BigDecimal dmxBalance = safeBD(dmx.get().getBalance()); + log.debug("OM*T.id={} -> DM*X.id={}, dmxBalance={}, omtBalance={}", + omtRgs.getId(), + dmx.get().getId(), + dmxBalance, omtBalance + ); + if (dmxBalance.compareTo(omtBalance) >= 0) { + group.forEach(rgs -> updateStatus(rgs, RegistryStatus.OK)); +// updateStatus(dmx.get(), RegistryStatus.POOL); + } else { + //для OM*T uncovered остальным фейл + group.forEach(rgs -> updateStatus(rgs, registryStatusFailed(rgs))); + updateStatus(dmx.get(), RegistryStatus.UNCV); + } + continue; + } + + Optional amfAssetO = searchAssetByOMT(omtRgs); + Optional amfAssetReceiverO = searchAssetByOMT(tmtRgs); + 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, registryStatusFailed())); // FAIL or MNG + continue; + } + if (amfAssetReceiverO.isEmpty()) { + log.error("TM*T register.id={} groupId={} failed to find AM*F asset", tmtRgs.getId(), entry.getKey()); + group.forEach(rgs -> updateStatus(rgs, registryStatusFailed())); // FAIL or MNG + continue; + } + + Registry amfAsset = amfAssetO.get(); + log.debug("groupId={}, AM*F.id={}", entry.getKey(), amfAsset.getId()); + BigDecimal amfBalance = safeBD(amfAsset.getBalance()); + BiConsumer processAssetsAndSetSessionId = (amf, amount) -> { + Optional asts = assets.processByAm_f(amf, amount); + if (asts.isPresent()) { + asts.get().a__b().setSessionId(sessionId); + registryImdg.update(asts.get().a__b()); + } + }; + Optional dmtInfo = rgsMng.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)); + processAssetsAndSetSessionId.accept(amfAsset, omtBalance.negate()); + processAssetsAndSetSessionId.accept(amfAssetReceiverO.get(), omtBalance); +// updateStatus(dmtInfo.get(), RegistryStatus.POOL); + continue; + } else if (SessionType.XDEP.equals(sessionType)) { + group.forEach(rgs -> updateStatus(rgs, registryStatusFailed(rgs))); // для финальной UNCV or FAIL ; другие MNG +// updateStatus(dmtInfo.get(), RegistryStatus.UNCV); не меняет статус, статус определяется OM*T + continue; + } + } + Optional dmtClnr = rgsMng.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)); + processAssetsAndSetSessionId.accept(amfAsset, omtBalance.negate()); + processAssetsAndSetSessionId.accept(amfAssetReceiverO.get(), omtBalance); + continue; + } else if (SessionType.XDEP.equals(sessionType)) { + group.forEach(rgs -> updateStatus(rgs, registryStatusFailed(rgs))); +// updateStatus(dmtClnr.get(), RegistryStatus.UNCV); не меняет статус, статус определяется OM*T + continue; + } + } + + if ((SessionType.FINL.equals(sessionType) || SessionType.UNIT.equals(sessionType)) + && amfBalance.compareTo(omtBalance) >= 0 + && initiatorV(omtRgs) && gatewayIfNeeded(omtRgs)) { + log.debug("OM*T.id={} -> no DM*X/DM*T(INFO/CLRN) registry found. " + + "FINL, initiator.category='V' and am*f > om*t: setting OK status to group", + omtRgs.getId()); + group.forEach(rgs -> updateStatus(rgs, RegistryStatus.OK)); + processAssetsAndSetSessionId.accept(amfAsset, omtBalance.negate()); + processAssetsAndSetSessionId.accept(amfAssetReceiverO.get(), omtBalance); + } else { + log.info("OM*T.id={} -> no DM*X/DM*T(INFO/CLRN) registry found. setting {} status to group", + omtRgs.getId(), registryStatusFailed()); + group.forEach(rgs -> updateStatus(rgs, registryStatusFailed(rgs))); + } + } + planBalanceCalc.stageRevision2(sessionId); + } + + 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()) + ); + } + + 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 RegistryStatus registryStatusFailed() { + if (SessionType.FINL.equals(sessionType) || SessionType.UNIT.equals(sessionType)) { + return RegistryStatus.FAIL; + } else { + return RegistryStatus.MNG; + } + } + + private RegistryStatus registryStatusFailed(Registry registry) { + if (SessionType.FINL.equals(sessionType) || SessionType.UNIT.equals(sessionType)) { + if (equalByRgs(OM_T, registry)) { + return RegistryStatus.UNCV; + } else { + return RegistryStatus.FAIL; + } + } else { + return RegistryStatus.MNG; + } + } + + private Optional searchAssetByOMT(Registry obligation) { + String sqlCondition = String.format("%s and " + + "tradingClearingRegistryId = '%s' and " + + "companyId = '%s' and " + + "securitySymbol = '%s'", + RegistryCodeSqlBuilder.getInstance(AM_F).build(), + obligation.getTradingClearingRegistryId(), + obligation.getCompanyId(), + obligation.getSecuritySymbol()); + Registry amf = registryImdg.getFirstObjectBySQL(sqlCondition); + return Optional.ofNullable(amf); + } + + private boolean gatewayIfNeeded(Registry omt) { + if (SessionType.FINL.equals(sessionType) && tradingTimeService.isTradingTime()) { + Optional gatewayReceived = gateway.gatewayRequestAndWait(() -> GatewayRequestCreator.from(omt)); + if (gatewayReceived.isEmpty()) { + log.error("gateway not received response for om*t.id: {}. ", omt.getId()); + } else { + log.debug("gateway response for om*t.id {}: approved {}", omt.getId(), gatewayReceived.get()); + } + return gatewayReceived.orElse(false); + } else { + log.trace("gateway not needed for om*t.id: {}", omt.getId()); + return true; + } + } + + private boolean initiatorV(Registry omtRgs) { + if (omtRgs.getCounterPartyId() == null) { + return false; + } + ClearingMemberCategory ctgr = categoryImdg.getFirstObjectBySQL( + "companyId = %d and clearingMemberCategory = '%s'" + .formatted(omtRgs.getCounterPartyId(), ClearingCategory.V.getKey())); + return ctgr != null; + } + + public void setSessionType(SessionType sessionType) { + this.sessionType = sessionType; + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/InspectionObligationsV2Action.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/InspectionObligationsV2Action.java new file mode 100644 index 000000000..652d5a2d6 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/InspectionObligationsV2Action.java @@ -0,0 +1,408 @@ +package ru.spcex.clearing.session.state.action; + +import java.math.BigDecimal; +import java.time.Instant; +import java.time.LocalDate; +import java.util.ArrayList; +import java.util.Collection; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.Optional; +import java.util.stream.Collectors; +import java.util.stream.Stream; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.config.ConfigurableBeanFactory; +import org.springframework.context.annotation.Scope; +import org.springframework.statemachine.StateContext; +import org.springframework.stereotype.Service; +import ru.clearing.classes.statics.data.registry.Registry; +import ru.spcex.clearing.component.predicate.cash.registry.AssetByObligationCashedPredicate; +import ru.spcex.clearing.component.predicate.cash.registry.AssetByTcrCompanyAccountCashedPredicate; +import ru.spcex.clearing.error.ClearingError; +import ru.spcex.clearing.error.RgsError; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.service.cash.impl.RegistryCashA__F; +import ru.spcex.clearing.service.cash.impl.RegistryCashAssetCash; +import ru.spcex.clearing.service.cash.impl.RegistryCashD___; +import ru.spcex.clearing.service.integration.GatewayRequestCreator; +import ru.spcex.clearing.service.registry.AssetTBFProcessing; +import ru.spcex.clearing.service.registry.AssetTBFProcessingCashing; +import ru.spcex.clearing.service.registry.RegistryManager; +import ru.spcex.clearing.service.schedule.TradingTimeService; +import ru.spcex.clearing.session.stage.TaskType; +import ru.spcex.clearing.session.stage.impl.GatewayRequester; +import ru.spcex.clearing.session.stage.impl.PlanBalanceCalc; +import ru.spcex.clearing.session.stage.util.RegistryUtil; +import ru.spcex.clearing.session.state.DataEnum; +import ru.spcex.clearing.session.state.SsnEvent; +import ru.spcex.platform.enumeration.RegistryDesignation; +import ru.spcex.platform.enumeration.RegistryInstrumentType; +import ru.spcex.platform.enumeration.RegistryStatus; +import static ru.spcex.platform.enumeration.RegistryTradingParams.*; +import ru.spcex.platform.enumeration.RegistryUnit; +import ru.spcex.platform.enumeration.Section; +import ru.spcex.platform.enumeration.SessionType; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgId; +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.collection.Pair; +import ru.spcex.platform.utils.enumeration.EnumMessage; +import ru.spcex.platform.utils.enumeration.IEnumKey; +import ru.spcex.platform.utils.enumeration.IMessageResolver; +import ru.spcex.platform.utils.enumeration.SimpleMessageResolver; +import static ru.spcex.platform.utils.number.BigDecimalUtil.safeBD; + +@Service +@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) +public class InspectionObligationsV2Action extends AbstractSessionActionForOkErrorHandling { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final Imdg registryImdg; + private final ImdgId idGen; + private final RegistryManager registryManager; + private final IMessageResolver msgResolver = new SimpleMessageResolver(); + private SessionType sessionType; + private Section section; + private final AssetTBFProcessing assets; + private final AssetTBFProcessingCashing assetsCashing; + private final GatewayRequester gateway; + private final TradingTimeService tradingTimeService; + private final PlanBalanceCalc planBalanceCalc; + + @Autowired + public InspectionObligationsV2Action(ImdgProvider imdgProvider, + RegistryManager registryManager, + AssetTBFProcessing assets, + GatewayRequester gateway, + PlanBalanceCalc planBalanceCalc + ) { + this.registryImdg = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_Registry, Registry.class, null); + this.idGen = imdgProvider.getImdgIdGenerator(); + this.registryManager = registryManager; + this.assets = assets; + this.gateway = gateway; + this.planBalanceCalc = planBalanceCalc; + this.gateway.setName("InspectionObligations|OM*T"); + this.tradingTimeService = new TradingTimeService(imdgProvider); + this.assetsCashing = new AssetTBFProcessingCashing(imdgProvider, d___Cash); + + } + + private final RegistryCashAssetCash a___Cash = new RegistryCashAssetCash("[all assets]"); + private final RegistryCashA__F a__fCash = new RegistryCashA__F("[A__F only]"); + private final RegistryCashD___ d___Cash = new RegistryCashD___("[D__I/V]"); + + + public void setSessionType(SessionType sessionType) { + this.sessionType = sessionType; + } + + + @Override + protected void actualExecute(StateContext ctx) { + try { + Long sessionId = ctx.getExtendedState().get(DataEnum.sessionId, Long.class); + inspectionObligations(sessionId); + } finally { + a___Cash.clear(); + a__fCash.clear(); + d___Cash.clear(); + } + } + + private void inspectionObligations(Long sessionId) { + String sqlCondition = String.format("(%s) and registryStatus = '%s'", + RegistryCodeSqlBuilder.getInstance(OS_T, OM_T, TS_T, TM_T).build(), + RegistryStatus.POOL.getKey()); + + Collection registriesToProcess = registryImdg.getCollectionObjectsBySQL(sqlCondition); + Map rgssToStore = new HashMap<>(registriesToProcess.size()); //вторые ноги у сделок - нужна ли оптимизация? + List>> registriesByGroupSorted = registriesToProcess.stream() + .collect(Collectors.groupingBy(Registry::getGroupId)) + .entrySet() + .stream() + .sorted((entry1, entry2) -> { + long minId1 = entry1.getValue().stream() + .mapToLong(Registry::getId) + .min() + .orElse(Long.MIN_VALUE); + + long minId2 = entry2.getValue().stream() + .mapToLong(Registry::getId) + .min() + .orElse(Long.MIN_VALUE); + + return Long.compare(minId1, minId2); + }) + .toList(); + + log.info("found {} ({} groups) registries by sql: {}", registriesToProcess.size(), registriesByGroupSorted.size(), sqlCondition); + ImdgPredicateBuilder rgsPb = registryImdg.predicateBuilder(); + GROUP: + for (Map.Entry> entry : registriesByGroupSorted) { + List group = entry.getValue(); + + List obligationsInGroup = group.stream().filter(registry -> + IEnumKey.getEnumByKey(RegistryDesignation.class, registry.getRegistryDesignation()) == RegistryDesignation.O).toList(); + log.debug("Find {} obligation with ids: {} in group: {}", obligationsInGroup.size(), + obligationsInGroup.stream() + .map(Registry::getId).collect(Collectors.toList()), + entry.getKey()); + 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())) { + checkResults.add(new CheckResult(obligation, false)); + continue; + } + ImdgPredicate assetPrdct = AssetByObligationCashedPredicate + .getPredicate(obligation) + .cashed(rgsPb, a__fCash); + Registry asset = registryImdg.getFirstObjectByPredicate(assetPrdct); + if (asset == null) { + log.warn("Not found asset by registry.id={}", obligation.getId()); + isUncovered = true; + } else if (obligation.getBalance().compareTo(asset.getBalance()) > 0) { + log.warn("groupId: {}: {}", obligation.getGroupId(), + msgResolver.resolve(new EnumMessage(ClearingError.InsecurityObligation, obligation.getCompanyId()))); + isUncovered = true; + } + checkResults.add(new CheckResult(obligation, isUncovered)); + isUncovered = false; + } + + Optional omt = group.stream().filter(rgs -> RegistryManager.equalsByCode(OM_T, rgs)).findFirst(); + for (Registry rgs : group) { + Runnable failGroup = () -> Stream.concat(group.stream(), getRefundDateRgsIfPresent(group).stream()) + .forEach(registry -> { + registry.setComment(msgResolver.resolve(RgsError.A__tNotFound)); + updateRegistryStatus(registry, RegistryStatus.FAIL, rgssToStore); + }); + RegistryInstrumentType mOrS = IEnumKey.getEnumByKey(RegistryInstrumentType.class, rgs.getRegistryInstrumentType()); + if (mOrS == null) { + log.error("RegistryInstrumentType is null for groupId {} {}.id={}", entry.getKey(), + rgs.getRegistryCode(), rgs.getId()); + failGroup.run(); + continue; + } + if (Section.MKR.equalsByKey(rgs.getSection()) && mOrS.equals(RegistryInstrumentType.S)) { + continue; + } + ImdgPredicate assetPredicate = AssetByTcrCompanyAccountCashedPredicate + .getPredicate(rgs, A__T.type(mOrS)) + .cashed(rgsPb, a___Cash); + Registry a__t = registryImdg.getFirstObjectByPredicate(assetPredicate); + if (a__t == null) { + log.debug("groupId {}, {}.id={} - A**T not found. settings FAIL to group", + entry.getKey(), rgs.getRegistryCode(), rgs.getId()); + failGroup.run(); + continue GROUP; + } + } + if (SessionType.FINL.equals(sessionType) + && omt.isPresent() + && checkResults.stream().noneMatch(checkResult -> checkResult.isUncovered) + && tradingTimeService.isTradingTime()) { + Optional gatewayReceived = gateway.gatewayRequestAndWait(() -> GatewayRequestCreator.from(omt.get())); + if (gatewayReceived.isEmpty()) { + log.error("Gateway not received response for groupId: {}. ", entry.getKey()); + } + if (!gatewayReceived.orElse(false)) { + checkResults.stream() + .filter(chk -> chk.registry().getId().equals(omt.get().getId())) + .findFirst() + .ifPresent(chk -> chk.isUncovered = true); + } + } + defineStatusAndUpdateRegistry(checkResults, group, rgssToStore); + if (checkResults.stream().noneMatch(checkResult -> checkResult.isUncovered)) { + Instant now = Instant.now(); + //по каждому регистру OM*T, TM*T, OS*T, TS*T из одной группы + for (Registry registry : group) {//TODO CLRNWORM проверка на пустой blockedRegistry + { + RegistryInstrumentType mOrS = IEnumKey.getEnumByKey(RegistryInstrumentType.class, registry.getRegistryInstrumentType()); + if (mOrS == null) continue; + ImdgPredicate a__tPrdct = AssetByTcrCompanyAccountCashedPredicate + .getPredicate(registry, A__T.type(mOrS)) + .cashed(rgsPb, a___Cash); + Registry a__t = registryImdg.getFirstObjectByPredicate(a__tPrdct); + Registry a__b = null; + Registry a__f = null; + if (a__t != null) { + ImdgPredicate a__bPrdct = AssetByTcrCompanyAccountCashedPredicate + .getPredicate(registry, A__B.type(mOrS)) + .cashed(rgsPb, a___Cash); + ImdgPredicate a__fPrdct = AssetByTcrCompanyAccountCashedPredicate + .getPredicate(registry, A__F.type(mOrS)) + .cashed(rgsPb, a___Cash); + + a__b = Optional.ofNullable( + registryImdg.getFirstObjectByPredicate(a__bPrdct) + ).orElseGet(() -> copyB(a__t)); + + a__f = Optional.ofNullable( + registryImdg.getFirstObjectByPredicate(a__fPrdct) + ).orElseGet(() -> copyF(a__t)); + + } + { + String msg; + if (a__t != null) { + msg = "A**T.id=%d/A**B.id=%d/A**F.id=%d".formatted( + a__t.getId(), a__b.getId(), a__f.getId() + ); + } else { + msg = "A**T/A**B/A**F not found"; + } + log.debug("registry {}.id={}: {}", registry.getRegistryCode(), registry.getId(), msg); + } + + if (a__t != null) { + BigDecimal amount = safeBD(registry.getBalance()); + if (IEnumKey.getEnumByKey(RegistryDesignation.class, registry.getRegistryDesignation()) == RegistryDesignation.O) { + amount = amount.negate(); + } else if (IEnumKey.getEnumByKey(RegistryDesignation.class, registry.getRegistryDesignation()) == RegistryDesignation.T) { + amount = amount; + } + a__b.setSessionId(sessionId); + assetsCashing.process(a__b, a__t, a__f, amount); + rgssToStore.put(a__b.getId(), a__b); + rgssToStore.put(a__f.getId(), a__f); + } + } + } + } + } + registryImdg.putAll(rgssToStore, 200); + + planBalanceCalc.stageRevision2(sessionId); + } + + + private Registry copyF(Registry rgs) { + Registry rgsF = rgs.clone(); + RegistryManager.zeroState(rgsF); + rgsF.setRegistryUnit(RegistryUnit.F.getKey()); + rgsF.setRegistryCode(RegistryUtil.clearingCode(rgsF)); + rgsF.setId(idGen.nextId()); + registryImdg.insert(rgsF); + log.debug("created {}.id={} by {}.id={}", rgsF.getRegistryCode(), rgsF.getId(), rgs.getRegistryCode(), rgs.getId()); + return rgsF; + } + + private Registry copyB(Registry rgs) { + Registry rgsB = rgs.clone(); + RegistryManager.zeroState(rgsB); + rgsB.setRegistryUnit(RegistryUnit.B.getKey()); + rgsB.setRegistryCode(RegistryUtil.clearingCode(rgsB)); + rgsB.setId(idGen.nextId()); + registryImdg.insert(rgsB); + log.debug("created {}.id={} by {}.id={}", rgsB.getRegistryCode(), rgsB.getId(), rgs.getRegistryCode(), rgs.getId()); + return rgsB; + } + + private void defineStatusAndUpdateRegistry(List checkResults, List registries, Map rgssToStore) { + boolean isOneUncovered = checkResults.stream().anyMatch(checkResult -> checkResult.isUncovered); + log.debug("groupId {}, uncovered: {}", registries.stream().findFirst().map(Registry::getGroupId).orElse(null), isOneUncovered); + String commentErr = msgResolver.resolve(ClearingError.InsecurityObligation, checkResults.stream() + .filter(chk -> chk.isUncovered) + .map(chk -> chk.registry.getCompanyId()) + .findFirst().orElse(null)); + + if (isOneUncovered) { + for (CheckResult checkResult : checkResults) { + Optional tRegistryWithSameCompany = registries.stream().filter(registry -> + registry.getCompanyId().equals(checkResult.registry.getCompanyId()) && + !registry.getRegistryDesignation().equals(checkResult.registry.getRegistryDesignation()) + ).findFirst(); + checkResult.registry.setComment(commentErr); + if (checkResult.isUncovered) { + updateRegistryStatus(checkResult.registry, uncvStatus(), rgssToStore); + tRegistryWithSameCompany.ifPresent(registry -> updateRegistryStatus(registry, failStatus(), rgssToStore)); + } else { + updateRegistryStatus(checkResult.registry, failStatus(), rgssToStore); + tRegistryWithSameCompany.ifPresent(registry -> updateRegistryStatus(registry, failStatus(), rgssToStore)); + } + } + //проверяем есть ли второй день для сделки (он не входит в пул, поэтому ищем отдельно) + getRefundDateRgsIfPresent(registries).forEach(rgs -> updateRegistryStatus(rgs, RegistryStatus.FAIL, rgssToStore)); + } else { + registries.forEach(registry -> updateRegistryStatus(registry, RegistryStatus.OK, rgssToStore)); + } + } + + private Collection getRefundDateRgsIfPresent(List registries) { + if (registries.stream().anyMatch(rgs -> rgs.getRefundDate() != null) && !SessionType.MEDM.equals(sessionType)) { + Optional> groupIdAndSettleDate = registries.stream() + .map(rgs -> new Pair<>(rgs.getGroupId(), rgs.getSettlementDate())) + .filter(pair -> pair.getFirst() != null) + .filter(pair -> pair.getSecond() != null) + .findFirst(); + if (groupIdAndSettleDate.isEmpty()) { + return Collections.emptyList(); + } + ImdgPredicateBuilder pb = registryImdg.predicateBuilder(); + return registryImdg.getCollectionObjectsByPredicate( + pb.and( + pb.equals("groupId", groupIdAndSettleDate.get().getFirst()), + pb.not(pb.equals("settlementDate", groupIdAndSettleDate.get().getSecond())) + ) + ); + } else return Collections.emptyList(); + } + + private void updateRegistryStatus(Registry registry, RegistryStatus registryStatus) { + log.debug("Update registry.id: {} to {}", registry.getId(), registryStatus.getKey()); + registry.setRegistryStatus(registryStatus.getKey()); + registry.setUpdated(Instant.now()); + registryImdg.update(registry); + } + + private void updateRegistryStatus(Registry registry, RegistryStatus registryStatus, Map rgss) { + log.debug("Update registry.id: {} to {}", registry.getId(), registryStatus.getKey()); + registry.setRegistryStatus(registryStatus.getKey()); + registry.setUpdated(Instant.now()); + rgss.put(registry.getId(), registry); + //registryImdg.update(registry); + } + + private RegistryStatus uncvStatus() { + if (SessionType.MEDM.equals(sessionType)) { + return RegistryStatus.MNG; + } else { + return RegistryStatus.UNCV; + } + } + + private RegistryStatus failStatus() { + if (SessionType.MEDM.equals(sessionType)) { + return RegistryStatus.MNG; + } else { + return RegistryStatus.FAIL; + } + } + + private static final class CheckResult { + private final Registry registry; + private boolean isUncovered; + + private CheckResult(Registry registry, boolean isUncovered) { + this.registry = registry; + this.isUncovered = isUncovered; + } + + public Registry registry() { + return registry; + } + + public boolean isUncovered() { + return isUncovered; + } + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/ObligationAdmissionAction.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/ObligationAdmissionAction.java new file mode 100644 index 000000000..3869e5484 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/ObligationAdmissionAction.java @@ -0,0 +1,135 @@ +package ru.spcex.clearing.session.state.action; + +import java.util.ArrayList; +import java.util.Collection; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.Optional; +import java.util.stream.Collectors; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.config.ConfigurableBeanFactory; +import org.springframework.context.annotation.Scope; +import org.springframework.statemachine.StateContext; +import org.springframework.stereotype.Service; +import ru.clearing.classes.statics.data.account.Account; +import ru.clearing.classes.statics.data.company.Company; +import ru.clearing.classes.statics.data.company.relation.Relation; +import ru.clearing.classes.statics.data.registry.Registry; +import ru.clearing.classes.statics.data.registry.TradingClearingRegistry; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.service.validation.admission.AccountActiveValidationRule; +import ru.spcex.clearing.service.validation.admission.ClearingAvailableValidationRule; +import ru.spcex.clearing.service.validation.admission.CompanyActiveValidationRule; +import ru.spcex.clearing.service.validation.admission.TcrActiveValidationRule; +import ru.spcex.clearing.session.stage.TaskType; +import ru.spcex.clearing.session.state.DataEnum; +import ru.spcex.clearing.session.state.SsnEvent; +import ru.spcex.platform.enumeration.RegistryStatus; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; +import ru.spcex.platform.imdg.iml.hazelcast.adapter.CashCloser; +import ru.spcex.platform.imdg.iml.hazelcast.adapter.CashV2ById; +import ru.spcex.platform.imdg.iml.hazelcast.adapter.CashV2ByIdAndString; +import ru.spcex.platform.imdg.validation.ImdgValidationContext; +import ru.spcex.platform.imdg.validation.LogPrefixId; +import ru.spcex.platform.utils.enumeration.EnumMessage; +import ru.spcex.platform.utils.enumeration.IMessageResolver; +import ru.spcex.platform.utils.validation.IValidator; +import ru.spcex.platform.utils.validation.ValidatorImpl; + +@Service +@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) +public class ObligationAdmissionAction extends AbstractSessionActionForOkErrorHandling { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final Imdg registryImdg; + private final Imdg imdgRltn; + private final Imdg imdgAcc; + private final Imdg imdgCmp; + private final Imdg imdgTcr; + private final IMessageResolver messageResolver; + private final ObligationAdmissionCash cash = new ObligationAdmissionCash(); + + + @Autowired + public ObligationAdmissionAction(ImdgProvider imdgProvider, IMessageResolver messageResolver) { + this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class); + this.messageResolver = messageResolver; + this.imdgRltn = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_Relation, Relation.class, null); + this.imdgAcc = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_Account, Account.class, cash.accCash); + this.imdgCmp = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_Company, Company.class, cash.cmpCash); + this.imdgTcr = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class, cash.tcrCash); + } + + @Override + protected void actualExecute(StateContext ctx) { + try (cash) { + Long sessionId = ctx.getExtendedState().get(DataEnum.sessionId, Long.class); + obligationAdmission(sessionId); + } + } + + private void obligationAdmission(Long sessionId) { + Collection registries = registryImdg.getCollectionObjectsBySQL("sessionId = " + sessionId); + Map> byGroups = registries.stream().collect(Collectors.groupingBy(Registry::getGroupId)); + log.info("found {} ({} groups) registries with sessionId {}", registries.size(), byGroups.size(), sessionId); + Map rgsToUpdate = new HashMap<>(); + for (Map.Entry> grpEntry : byGroups.entrySet()) { + EnumMessage groupError = null; + List rgsGroup = grpEntry.getValue(); + for (Registry rgs : rgsGroup) { + IValidator validator = valFor(rgs); + Optional error = validator.tillFirstError(); + if (error.isPresent()) { + groupError = error.get(); + break; + } + } + if (groupError != null) { + log.warn("error {} for registries groupId = {}", messageResolver.resolve(groupError), grpEntry.getKey()); + for (Registry rgs : rgsGroup) { + rgs.setRegistryStatus(RegistryStatus.NACK.getKey()); + rgsToUpdate.put(rgs.getId(), rgs); + } + } + } + registryImdg.putAll(rgsToUpdate); + } + + private IValidator valFor(Registry rgs) { + ImdgValidationContext ctx = new ImdgValidationContext<>(); + ctx.setValidatedObject(rgs); + ctx.addImdg(IMDGDistributedNames.Map_Relation, imdgRltn); + //ctx.addImdg(IMDGDistributedNames.Map_Session, imdgSession); + //ctx.addImdg(IMDGDistributedNames.Map_SectionDictionary, imdgSectionDictionary); + ctx.addImdg(IMDGDistributedNames.Map_Account, imdgAcc); + ctx.addImdg(IMDGDistributedNames.Map_Company, imdgCmp); + ctx.addImdg(IMDGDistributedNames.Map_TradingClearingRegistry, imdgTcr); + ctx.setLogPrefix(LogPrefixId.INSTANCE); + return new ValidatorImpl<>(ctx, + new ClearingAvailableValidationRule(cash.rltnCash), + new AccountActiveValidationRule(), + new CompanyActiveValidationRule(), + new TcrActiveValidationRule()); + } + + + protected static class ObligationAdmissionCash extends CashCloser { + CashV2ByIdAndString rltnCash; + CashV2ById accCash; + CashV2ById cmpCash; + CashV2ById tcrCash; + + public ObligationAdmissionCash() { + this.cashes = new ArrayList<>(); + rltnCash = add(new CashV2ByIdAndString<>("rltnCash", rltn -> + new CashV2ByIdAndString.CustomKey(rltn.getConsumerId(), rltn.getService()))); + accCash = add(new CashV2ById<>("accCash")); + cmpCash = add(new CashV2ById<>("cmpCash")); + tcrCash = add(new CashV2ById<>("tcrCash")); + } + } + +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/RequirementAndObligationCreationAction.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/RequirementAndObligationCreationAction.java new file mode 100644 index 000000000..ee18097e8 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/RequirementAndObligationCreationAction.java @@ -0,0 +1,454 @@ +package ru.spcex.clearing.session.state.action; + +import java.time.LocalDate; +import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.Objects; +import java.util.Optional; +import java.util.concurrent.atomic.AtomicReference; +import java.util.function.BiFunction; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.config.ConfigurableBeanFactory; +import org.springframework.context.annotation.Scope; +import org.springframework.statemachine.StateContext; +import org.springframework.stereotype.Service; +import ru.clearing.classes.statics.data.account.Account; +import ru.clearing.classes.statics.data.account.ClearingAccount; +import ru.clearing.classes.statics.data.account.DepoAccount; +import ru.clearing.classes.statics.data.company.Company; +import ru.clearing.classes.statics.data.execution.ExecutionCommon; +import ru.clearing.classes.statics.data.execution.ExecutionCurrency; +import ru.clearing.classes.statics.data.execution.ExecutionDeposit; +import ru.clearing.classes.statics.data.execution.ExecutionFond; +import ru.clearing.classes.statics.data.instrument.issue.EquitySecurity; +import ru.clearing.classes.statics.data.instrument.issue.FixedIncomeSecurity; +import ru.clearing.classes.statics.data.misc.Currency; +import ru.clearing.classes.statics.data.misc.Session; +import ru.clearing.classes.statics.data.registry.Registry; +import ru.clearing.classes.statics.data.registry.TradingClearingRegistry; +import ru.clearing.classes.statics.data.registry.TradingClearingRegistryList; +import ru.clearing.classes.statics.data.security.CurrencyPairSecurity; +import ru.clearing.classes.statics.data.security.MoneyMarketSecurity; +import ru.clearing.classes.statics.data.security.Security; +import ru.clearing.platform.dictionary.CurrencyPairDictionary; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.service.TcrSearcher; +import ru.spcex.clearing.session.stage.TaskType; +import ru.spcex.clearing.session.stage.impl.IRegistryBuilder; +import ru.spcex.clearing.session.stage.impl.RegistryCurrBuilder; +import ru.spcex.clearing.session.stage.impl.RegistryDepoBuilder; +import ru.spcex.clearing.session.stage.impl.RegistryFondBuilder; +import ru.spcex.clearing.session.stage.impl.RegistryUpdater; +import ru.spcex.clearing.session.state.DataEnum; +import ru.spcex.clearing.session.state.SsnEvent; +import static ru.spcex.clearing.util.ComparatorUtil.execIdComparator; +import static ru.spcex.clearing.util.ComparatorUtil.tradeTimeComparator; +import ru.spcex.platform.classes.base.interfaces.ExecutionType; +import ru.spcex.platform.enumeration.ISide; +import ru.spcex.platform.enumeration.MoneyFlowSide; +import ru.spcex.platform.enumeration.RegistryCapacity; +import ru.spcex.platform.enumeration.RegistryDesignation; +import ru.spcex.platform.enumeration.RegistryInstrumentType; +import ru.spcex.platform.enumeration.RegistryTradingParams; +import ru.spcex.platform.enumeration.RegistryUnit; +import ru.spcex.platform.enumeration.Side; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgId; +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.imdg.api.predicate.specific.SecuritySelector; +import ru.spcex.platform.imdg.iml.hazelcast.adapter.CashCloser; +import ru.spcex.platform.imdg.iml.hazelcast.adapter.CashV2; +import ru.spcex.platform.imdg.iml.hazelcast.adapter.CashV2ById; +import ru.spcex.platform.imdg.iml.hazelcast.adapter.CashV2ByIdAndString; +import ru.spcex.platform.imdg.iml.hazelcast.adapter.CashV2ByString; +import ru.spcex.platform.utils.collection.Pair; + +@Service +@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) +public class RequirementAndObligationCreationAction extends AbstractSessionActionForOkErrorHandling { + private final Logger log = LoggerFactory.getLogger(getClass()); + protected final Imdg registryImdg; + private final ImdgId idGen; + private final ImdgProvider imdgProvider; + protected final RequirementsAndObligationCreationCash cash; + + private final Imdg cmpImdg; + private final Imdg tcrImdg; + private final Imdg accImdg; + private final Imdg ssnImdg; + private final Imdg clAccImdg; + private final Imdg dpAccImdg; + private final Imdg currPairSecImdg; + private final Imdg fixIncSecImdg; + private final Imdg monMrktSecImdg; + private final Imdg eqtySecImdg; + private final Imdg currPairDictImdg; + private final Imdg currImdg; + private final TcrSearcher tcrSearcher; + private final SecuritySelector secSelector; + + @Autowired + public RequirementAndObligationCreationAction(ImdgProvider imdgProvider) { + this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class); + this.idGen = imdgProvider.getImdgIdGenerator(); + this.imdgProvider = imdgProvider; + this.cash = new RequirementsAndObligationCreationCash(); + this.cmpImdg = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_Company, Company.class, cash.cmpCash); + this.tcrImdg = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class, cash.tcrCash); + Imdg tcrListImdg = imdgProvider.getCashingImdg( + IMDGDistributedNames.Map_TradingClearingRegistryList, TradingClearingRegistryList.class, null + ); + this.accImdg = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_Account, Account.class, cash.accCash); + this.ssnImdg = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_Session, Session.class, cash.ssnCash); + this.clAccImdg = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_ClearingAccount, ClearingAccount.class, null); + this.dpAccImdg = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_DepoAccount, DepoAccount.class, null); + this.currPairSecImdg = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_CurrencyPairSecurity, CurrencyPairSecurity.class, cash.currPairCash); + this.fixIncSecImdg = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_FixedIncomeSecurity, FixedIncomeSecurity.class, cash.fixIncSecCash); + this.monMrktSecImdg = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_MoneyMarketSecurity, MoneyMarketSecurity.class, cash.monMrktSecCash); + this.eqtySecImdg = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_EquitySecurity, EquitySecurity.class, cash.eqtySecCash); + this.currPairDictImdg = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_CurrencyPairDictionary, CurrencyPairDictionary.class, cash.currPairDictCash); + this.currImdg = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_Currency, Currency.class, null); + + + this.tcrSearcher = new TcrSearcher(tcrImdg, tcrListImdg, cmpImdg, accImdg); + this.tcrSearcher.setTcrListCash(cash.tcrListCash); + this.secSelector = new SecuritySelector<>(castSecImdg(fixIncSecImdg), + castSecImdg(monMrktSecImdg), + castSecImdg(eqtySecImdg), + castSecImdg(currPairSecImdg)); + } + + private static Imdg castSecImdg(Imdg m) { + return (Imdg) m; + } + + + protected Boolean doRgsSearch = null; + + @Override + public void actualExecute(StateContext ctx) { + try (cash) { + @SuppressWarnings("unchecked") + List list = ctx.getExtendedState().get(DataEnum.dealsPrepared, List.class); + Long sessionId = ctx.getExtendedState().get(DataEnum.sessionId, Long.class); + createRegisters(list, sessionId, null); + } finally { + doRgsSearch = null; + } + } + + protected void createRegisters(List data, Long sessionId, Map rgsStorage) { + data.sort(tradeTimeComparator.thenComparing(execIdComparator)); + Map newRgss = rgsStorage != null ? rgsStorage : storageMap(data); + //см. описание к #matchExecutions + for (int i = 0; i < data.size(); ) { + Pair matched = matchExecutions(data, i); + if (matched == null) { + log.error("couldn't find match execution for id = {}", data.get(i).getId()); + i += 1; + continue; + } + log.debug("matched executions id {} and {}", matched.getFirst().getId(), matched.getSecond().getId()); + i += 2; + ExecutionCommon partyExec = matched.getFirst(); + ExecutionCommon counterExec = matched.getSecond(); + //бывает поиск обычный, бывает по ExecutionDeposit.secondLegSettlementDt + AtomicReference>> regSearchMethod = new AtomicReference<>( + (exc, dsg) -> Optional.empty() + ); + //обычный поиск + if (searchForRegisters(sessionId)) { + regSearchMethod.set(this::findReg); + } + TriConsumer findAndUpdateOrCreate = (exec, regDsgn, returnedDep) -> + regSearchMethod.get().apply(exec, regDsgn) + .ifPresentOrElse(registry -> { + RegistryUpdater.updater() + .exec(exec) + .registry(registry) + .update(); + registryImdg.update(registry); + log.debug("executions id {} and {}: updated rgs.id={}", partyExec.getId(), counterExec.getId(), registry.getId()); + }, () -> { + IRegistryBuilder registryBuilder = defineBuilderByExecType(partyExec.type()); + Registry newRegister = registryBuilder + //cashes + .dpAccCash(cash.dpAccCash) + .currCash(cash.currCash) + .clrAccCash(cash.clrAccCash) + //imdg + .cmpImdg(cmpImdg) + .tcrImdg(tcrImdg) + .accImdg(accImdg) + .ssnImdg(ssnImdg) + .clAccImdg(clAccImdg) + .dpAccImdg(dpAccImdg) + .currPairSecImdg(currPairSecImdg) + .currPairDictImdg(currPairDictImdg) + .currImdg(currImdg) + .tcrSearcher(tcrSearcher) + .secSelector(secSelector) + .imdg(imdgProvider) + .exec(exec) + .registryDesignation(regDsgn) + .returnDeposit(returnedDep) + .build(); + newRegister.setId(idGen.nextId()); + newRgss.put(newRegister.getId(), newRegister); + log.debug("executions id {} and {}: created rgs.id={}", partyExec.getId(), counterExec.getId(), newRegister.getId()); + }); + findAndUpdateOrCreate.accept(partyExec, RegistryDesignation.O, false); + findAndUpdateOrCreate.accept(partyExec, RegistryDesignation.T, false); + findAndUpdateOrCreate.accept(counterExec, RegistryDesignation.O, false); + findAndUpdateOrCreate.accept(counterExec, RegistryDesignation.T, false); + if (partyExec.type().equals(ExecutionType.ExecutionDeposit) && ((ExecutionDeposit) partyExec).getSecondLegSettlementDate() != null) { + //поиск по дате расчетов ExecutionDeposit.secondLegSettlementDt + if (searchForRegisters(sessionId)) { + regSearchMethod.set(this::findRegBySettlementDt); + } + findAndUpdateOrCreate.accept(partyExec, RegistryDesignation.T, true); + findAndUpdateOrCreate.accept(partyExec, RegistryDesignation.O, true); + findAndUpdateOrCreate.accept(counterExec, RegistryDesignation.T, true); + findAndUpdateOrCreate.accept(counterExec, RegistryDesignation.O, true); + } + } + if (rgsStorage == null) { + log.info("batch insert {} registries into map", newRgss.size()); + registryImdg.putAll(newRgss); + } + } + + + + protected boolean searchForRegisters(Long sessionId) { + if (sessionId == null) return true; + if (doRgsSearch != null) return doRgsSearch; + + ImdgPredicateBuilder pb = registryImdg.predicateBuilder(); + ImdgPredicate condition = pb.and( + pb.equals("sessionId", sessionId), + pb.or( + pb.equals("registryDesignation", RegistryDesignation.O.getKey()), + pb.equals("registryDesignation", RegistryDesignation.T.getKey()) + ), + pb.or( + pb.equals("registryInstrumentType", RegistryInstrumentType.M.getKey()), + pb.equals("registryInstrumentType", RegistryInstrumentType.S.getKey()) + ), + pb.equals("registryUnit", RegistryUnit.T.getKey()) + ); + Registry rgs = registryImdg.getFirstObjectByPredicate(condition); + this.doRgsSearch = rgs != null; + return this.doRgsSearch; + } + + protected static int calcMapCapacity(int num) { + double divided = num / 0.75; + return (int) Math.ceil(divided); + } + + protected Map storageMap(List data) { + ExecutionType eType = data.stream() + .findFirst() + .map(ExecutionCommon::type) + .orElse(ExecutionType.ExecutionDeposit); + int initialCapacity; + if (eType.equals(ExecutionType.ExecutionDeposit)) { + initialCapacity = calcMapCapacity(data.size() * 4); + } else { + initialCapacity = calcMapCapacity(data.size() * 2); + } + log.trace("executions size {}; storage map initial capacity {}", data.size(), initialCapacity); + return new HashMap<>(initialCapacity, 1f); + } + + private Optional findReg(ExecutionCommon exec, RegistryDesignation des) { + ISide side = getSide(exec); + RegistryTradingParams p; + if (side.isBuy() && des.equals(RegistryDesignation.O)) { + p = new RegistryTradingParams( + RegistryDesignation.O, RegistryInstrumentType.M, null, RegistryUnit.T + ); + } else if (side.isBuy() && des.equals(RegistryDesignation.T)) { + p = new RegistryTradingParams( + RegistryDesignation.T, RegistryInstrumentType.S, null, RegistryUnit.T + ); + } else if (side.isSell() && des.equals(RegistryDesignation.O)) { + p = new RegistryTradingParams( + RegistryDesignation.O, RegistryInstrumentType.S, null, RegistryUnit.T + ); + } else if (side.isSell() && des.equals(RegistryDesignation.T)) { + p = new RegistryTradingParams( + RegistryDesignation.T, RegistryInstrumentType.M, null, RegistryUnit.T + ); + } else { + throw new IllegalStateException("cannot construct for " + des + " " + side); + } + return imdgRegSearch(p, + exec.getTradingClearingRegistryId(), + exec.getCompanyId(), + settlementDt(exec), + exec.getSessionId(), + exec.getExchangeExecutionId() + ); + } + + private Optional findRegBySettlementDt(ExecutionCommon exec, RegistryDesignation des) { + if (!exec.type().equals(ExecutionType.ExecutionDeposit) || ((ExecutionDeposit) exec).getSecondLegSettlementDate() == null) + throw new IllegalStateException("cannot searchSettlementDt for " + exec.type()); + ISide side = getSide(exec); + RegistryTradingParams p; + if (side.isBuy() && des.equals(RegistryDesignation.T)) { + p = new RegistryTradingParams( + RegistryDesignation.T, RegistryInstrumentType.M, RegistryCapacity.A, RegistryUnit.T + ); + } else if (side.isBuy() && des.equals(RegistryDesignation.O)) { + p = new RegistryTradingParams( + RegistryDesignation.O, RegistryInstrumentType.S, RegistryCapacity.A, RegistryUnit.T + ); + } else if (side.isSell() && des.equals(RegistryDesignation.O)) { + p = new RegistryTradingParams( + RegistryDesignation.O, RegistryInstrumentType.M, RegistryCapacity.A, RegistryUnit.T + ); + } else if (side.isSell() && des.equals(RegistryDesignation.T)) { + p = new RegistryTradingParams( + RegistryDesignation.T, RegistryInstrumentType.S, RegistryCapacity.A, RegistryUnit.T + ); + } else { + throw new IllegalStateException("cannot construct by settleDt for " + des + " " + side); + } + return imdgRegSearch(p, + exec.getTradingClearingRegistryId(), + exec.getCompanyId(), + ((ExecutionDeposit) exec).getSecondLegSettlementDate(), + exec.getSessionId(), + exec.getExchangeExecutionId()); + } + + public Optional imdgRegSearch(RegistryTradingParams p, Long tcrId, Long companyId, LocalDate settlementDt, + Long sessionId, Long exchangeExecutionId) { + String sql = RegistryCodeSqlBuilder.getInstance(p).build(); + ImdgPredicateBuilder pb = registryImdg.predicateBuilder(); + ImdgPredicate rgstrPredicate = pb.and(pb.sql(sql), + pb.equals("tradingClearingRegistryId", tcrId), + pb.equals("companyId", companyId), + pb.equals("settlementDate", settlementDt), + pb.equals("sessionId", sessionId), + pb.equals("groupId", exchangeExecutionId) + ); + Registry rgs = registryImdg.getSingleObjectByPredicate(rgstrPredicate); + log.trace("rgs.id={} found by: {}", rgs == null ? "'not found'" : rgs.getId(), rgstrPredicate.toString()); + return Optional.ofNullable(rgs); + } + + private LocalDate settlementDt(ExecutionCommon exec) { + if (exec.type().equals(ExecutionType.ExecutionDeposit)) { + return ((ExecutionDeposit) exec).getFirstLegSettlementDate(); + } else if (exec.type().equals(ExecutionType.ExecutionFond)) { + return ((ExecutionFond) exec).getSettlementDate(); + } else if (exec.type().equals(ExecutionType.ExecutionCurrency)) { + return ((ExecutionCurrency) exec).getSettlementDate(); + } else { + throw new RuntimeException("Unknown Execution type: " + exec.type()); + } + } + + + /** + * входные данные отсортированы по exchangeExecutionId. + * сортировка выполнена на предыдущем шаге. (такое описание шагов в ТЗ) + * по идее, Execution с одинаковым exchangeExecutionId должно быть всего два. + * если это не так, передалать на коллекцию Long executionId в правильном порядке + * и для каждого искать мэтч отдельно в Imdg, с сохранением уже обработанных для избежания дублирования + */ + protected Pair matchExecutions(List data, int i) { + //if the last execution, then no match + if (i >= data.size() - 1) { + return null; + } + ExecutionCommon exec1 = data.get(i); + ExecutionCommon exec2 = data.get(i + 1); + //Встречная сделка контрагента выбирается из executionDeposit/Fond по условию: + //exchangeExecutionId=currentExecutionDeposit/Fond.exchangeExecutionId + //и [side=SELL (если currentExecutionDeposit/Fond.side=BUY) или side=BUY (если currentExecutionDeposit/Fond.side=SELL) по справочнику moneyFlowSide или справочнику side в зависимости от секции обрабатываемой сделки] + //и companyId=currentExecutionDeposit/Fond.counterPartyId + String invalid = null; + if (!Objects.equals(exec1.getExchangeExecutionId(), exec2.getExchangeExecutionId())) { + invalid = String.format("exchangeExecutionId %d and %d not equal", exec1.getExchangeExecutionId(), exec2.getExchangeExecutionId()); + } else if (Objects.equals(getSide(exec1), getSide(exec2))) { + invalid = String.format("side %s and %s equal, must be opposite", getSide(exec1), getSide(exec1)); + } else if (!Objects.equals(exec1.getCompanyId(), exec2.getCounterPartyId()) || !Objects.equals(exec1.getCounterPartyId(), exec2.getCompanyId())) { + invalid = String.format("Execution#id(%d)#companyId(%d)#counterPartyId(%d), Execution#id(%d)#companyId(%d)#counterPartyId(%d)", + exec1.getId(), exec1.getCompanyId(), exec1.getCounterPartyId(), exec2.getId(), exec2.getCompanyId(), exec2.getCounterPartyId()); + } + if (invalid != null) { + log.error("FATAL couldn't match Execution#id({}) with Execution#id({}) error: {}", exec1.getId(), exec2.getId(), invalid); + return null; + } + return new Pair<>(exec1, exec2); + } + + private static ISide getSide(ExecutionCommon exec) { + if (exec.type().equals(ExecutionType.ExecutionDeposit)) { + return ISide.parse(MoneyFlowSide.class, exec.getSide()); + } else if (exec.type().equals(ExecutionType.ExecutionFond) || exec.type().equals(ExecutionType.ExecutionCurrency)) { + return ISide.parse(Side.class, exec.getSide()); + } else { + throw new RuntimeException("Unknown Execution type: " + exec.type()); + } + } + + private IRegistryBuilder defineBuilderByExecType(ExecutionType executionType) { + return switch (executionType) { + case ExecutionFond -> RegistryFondBuilder.builder(); + case ExecutionDeposit -> RegistryDepoBuilder.builder(); + case ExecutionCurrency -> RegistryCurrBuilder.builder(); + }; + } + + @FunctionalInterface + private static interface TriConsumer { + void accept(T1 t1, T2 t2, T3 t3); + } + + protected static class RequirementsAndObligationCreationCash extends CashCloser { + protected final CashV2 cmpCash; + protected final CashV2 tcrCash; + protected final CashV2 accCash; + protected final CashV2 dpAccCash; + protected final CashV2 currCash; + protected final CashV2 clrAccCash; + protected final CashV2 ssnCash; + protected final CashV2 currPairCash; + protected final CashV2 fixIncSecCash; + protected final CashV2 monMrktSecCash; + protected final CashV2 eqtySecCash; + protected final CashV2 currPairDictCash; + protected final CashV2ByIdAndString tcrListCash; + + public RequirementsAndObligationCreationCash() { + cashes = new ArrayList<>(); + cmpCash = add(new CashV2ById<>("cmpCash")); + tcrCash = add(new CashV2ById<>("tcrCash")); + accCash = add(new CashV2ById<>("accCash")); + dpAccCash = add(new CashV2ById<>("dpAccCash")); + currCash = add(new CashV2ByString<>("currCash", Currency::getCurrencyCode)); + clrAccCash = add(new CashV2ById<>("clrAccCash")); + ssnCash = add(new CashV2ById<>("ssnCash")); + currPairCash = add(new CashV2ById<>("currPairCash")); + fixIncSecCash = add(new CashV2ById<>("fixIncSecCash")); + monMrktSecCash = add(new CashV2ById<>("monMrktSecCash")); + eqtySecCash = add(new CashV2ById<>("eqtySecCash")); + currPairDictCash = add(new CashV2ById<>("currPairDictCash")); + tcrListCash = add(new CashV2ByIdAndString<>("tcrListCash", tcrList -> new CashV2ByIdAndString.CustomKey(tcrList.getTradingClearingRegistryId(), tcrList.getCurrency()))); + } + } + +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/ReviseStage1.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/ReviseStage1.java new file mode 100644 index 000000000..1ae547dcc --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/ReviseStage1.java @@ -0,0 +1,55 @@ +package ru.spcex.clearing.session.state.action; + +import java.time.Instant; +import java.util.Collection; +import java.util.Objects; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.config.ConfigurableBeanFactory; +import org.springframework.context.annotation.Scope; +import org.springframework.statemachine.StateContext; +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.TaskType; +import ru.spcex.clearing.session.state.SsnEvent; +import ru.spcex.platform.enumeration.RegistryTradingParams; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; +import ru.spcex.platform.imdg.api.predicate.specific.RegistryCodeSqlBuilder; + +@Service +@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) +public class ReviseStage1 extends AbstractSessionActionForOkErrorHandling { + private final Logger log = LoggerFactory.getLogger(getClass()); + private Imdg registryImdg; + + @Autowired + public ReviseStage1(ImdgProvider imdgProvider) { + this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class); + } + + @Override + public void actualExecute(StateContext ctx) { + String sql = RegistryCodeSqlBuilder.getInstance(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); + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/Sdf56Action.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/Sdf56Action.java new file mode 100644 index 000000000..0d1b765d2 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/Sdf56Action.java @@ -0,0 +1,97 @@ +package ru.spcex.clearing.session.state.action; + +import java.time.Instant; +import java.time.LocalDate; +import java.time.LocalDateTime; +import java.time.LocalTime; +import java.time.temporal.ChronoUnit; +import java.util.Collection; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.config.ConfigurableBeanFactory; +import org.springframework.context.annotation.Scope; +import org.springframework.statemachine.StateContext; +import org.springframework.stereotype.Service; +import ru.clearing.classes.statics.data.misc.Currency; +import ru.clearing.classes.statics.data.sdf.SDf56; +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.service.sender.KafkaSender; +import ru.spcex.clearing.session.stage.TaskType; +import ru.spcex.clearing.session.state.ActionHeader; +import ru.spcex.clearing.session.state.SsnEvent; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgId; +import ru.spcex.platform.imdg.api.ImdgProvider; +import ru.spcex.platform.imdg.api.predicate.specific.StatementRevisePredicate; +import ru.spcex.platform.utils.time.TimeUtil; + +@Service +@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) +public class Sdf56Action extends AbstractSessionActionForOkErrorHandling { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final ImdgId idGenerator; + private final Imdg statementImdg; + private final Imdg currencyImdg; + private final Imdg sDf56Imdg; + private final KafkaSender kafkaSender; + + @Autowired + public Sdf56Action(ImdgProvider imdgProvider, KafkaSender kafkaSender) { + this.idGenerator = imdgProvider.getImdgIdGenerator(); + this.statementImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Statement, Statement.class); + this.currencyImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Currency, Currency.class); + this.sDf56Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf56, SDf56.class); + this.kafkaSender = kafkaSender; + } + + @Override + public void actualExecute(StateContext ctx) { + LocalTime fromTime = null; + if (ctx.getMessageHeader(ActionHeader.sdf56FromTime) instanceof LocalTime) { + fromTime = (LocalTime) ctx.getMessageHeader(ActionHeader.sdf56FromTime); + } + sendSdfs56(fromTime); + } + + private void sendSdfs56(LocalTime fromTime) { + Long millis = null; + if (fromTime != null) { + log.debug("sendSdf 56, parameter fromTime={}", fromTime); + Instant fromInstant = TimeUtil.localDateTimeToInstant(LocalDateTime.of(LocalDate.now(), fromTime)); + millis = fromInstant.toEpochMilli(); + } else { + Collection currencies = currencyImdg.projectSingleAttribute("id"); + Statement statement = statementImdg.aggregateByMax("created", StatementRevisePredicate.get(statementImdg, currencies)); + if (statement != null) { + log.debug("sendSdf 56, parameter fromTime not set, use max(statement.created)={}", statement.getCreated()); + millis = statement.getCreated().toEpochMilli(); + } else { + log.debug("sendSdf 56, parameter fromTime not set, and statement's not found."); + } + } + newSDf56(millis); + } + + + private void newSDf56(Long millis) { + log.debug("creating sdf56"); + SDf56 sDf56 = new SDf56(); + sDf56.setNumber(idGenerator.nextId().toString()); + Instant now = Instant.now(); + String startTime = String.valueOf(millis != null ? + millis : now.minus(1, ChronoUnit.DAYS).toEpochMilli()); + sDf56.setStart_datetime(startTime); + sDf56.setEnd_datetime(String.valueOf(now.toEpochMilli())); + sDf56.setGenerationTime(now); + sDf56.setGenerationId(idGenerator.nextId()); + sDf56Imdg.insert(sDf56); + SdfClearingRequest requestForExporter = new SdfClearingRequest(); + requestForExporter.setGroupId(sDf56.getGenerationId()); + kafkaSender.sendRequestToQueue(Consts.SDF56_PROCESS, requestForExporter); + log.debug("successfully processed, new id {}", sDf56.getId()); + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/factory/SessionStateMachineWrapper.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/factory/SessionStateMachineWrapper.java new file mode 100644 index 000000000..b04050076 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/factory/SessionStateMachineWrapper.java @@ -0,0 +1,110 @@ +package ru.spcex.clearing.session.state.factory; + +import java.util.AbstractMap; +import java.util.Map; +import java.util.NoSuchElementException; +import java.util.stream.Collectors; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.messaging.Message; +import org.springframework.statemachine.StateMachine; +import org.springframework.statemachine.config.StateMachineFactory; +import org.springframework.stereotype.Component; +import ru.clearing.classes.statics.data.misc.Session; +import ru.spcex.clearing.platform.messaging.domain.BaseRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest; +import ru.spcex.clearing.session.stage.TaskType; +import ru.spcex.clearing.session.state.DataEnum; +import ru.spcex.clearing.session.state.SsnEvent; +import ru.spcex.clearing.session.state.interceptor.SessionStatusChangingInterceptor; +import ru.spcex.clearing.session.state.listener.MachineStopListener; +import ru.spcex.platform.enumeration.Section; +import ru.spcex.platform.enumeration.SessionType; +import ru.spcex.platform.utils.enumeration.IEnumKey; + +@Component +public class SessionStateMachineWrapper { + private final Logger log = LoggerFactory.getLogger(getClass()); + + private final Map> factories; + private final SessionStatusChangingInterceptor statusChangingInterceptor; + private StateMachine currentSession; + + @Autowired + public SessionStateMachineWrapper( + Map> factories, + SessionStatusChangingInterceptor statusChangingInterceptor + ) { + this.factories = factories + .entrySet() + .stream() + .map(e -> { + SessionType sessType = IEnumKey.getEnumByKey(SessionType.class, e.getKey()); + if (sessType == null) throw new IllegalStateException("unknown session type " + e.getKey()); + return new AbstractMap.SimpleEntry<>(sessType, e.getValue()); + }) + .collect(Collectors.toMap(AbstractMap.SimpleEntry::getKey, AbstractMap.SimpleEntry::getValue)); + this.statusChangingInterceptor = statusChangingInterceptor; + } + + public synchronized void defineAndStartSession(BaseRequest r) { + LauncherCommandRequest payload = r.getRequestPayload(); + SessionType sessionType = IEnumKey.getEnumByKey(SessionType.class, payload.getSessionType()); + Section section = IEnumKey.getEnumByKey(Section.class, payload.getSection()); + if (sessionType == null || section == null) { + log.warn("Unknown sessionType: {} or section: {}", section, sessionType); + return; + } + StateMachineFactory factory = getByType(sessionType); + if (this.currentSession != null) { + Session session = currentSession + .getExtendedState() + .get(DataEnum.session, Session.class); + Long id = session != null ? session.getId() : null; + String ssnType = session != null ? session.getSessionType() : null; + log.warn("Session already running: type={} id={}", ssnType, id); + return; + } + this.currentSession = build(factory); + this.currentSession.addStateListener(new MachineStopListener( + this::clearCurrentMachine + )); + this.currentSession.start(); + } + + public synchronized void sendEvent(Message event) { + if (this.currentSession != null) { + log.info("sending an event to session: {}", event.getPayload()); + this.currentSession.sendEvent(event); + } else { + log.warn("there is no currently active session"); + } + } + + private synchronized void clearCurrentMachine() { + this.currentSession.stop(); + this.currentSession = null; + } + + private StateMachineFactory getByType(SessionType sessionType) { + StateMachineFactory factory = factories.get(sessionType); + if (factory == null) { + throw new NoSuchElementException("no state machine for session: " + sessionType); + } + return factory; + } + + + private StateMachine build( + StateMachineFactory stateMachineFactory + ) { + StateMachine stateMachine = stateMachineFactory.getStateMachine(); + stateMachine + .getStateMachineAccessor() + .doWithAllRegions( + access -> access.addStateMachineInterceptor(statusChangingInterceptor) + ); + return stateMachine; + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/interceptor/SessionStatusChangingInterceptor.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/interceptor/SessionStatusChangingInterceptor.java new file mode 100644 index 000000000..06386496d --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/interceptor/SessionStatusChangingInterceptor.java @@ -0,0 +1,108 @@ +package ru.spcex.clearing.session.state.interceptor; + +import java.time.Instant; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.statemachine.StateContext; +import org.springframework.statemachine.StateMachine; +import org.springframework.statemachine.state.State; +import org.springframework.statemachine.support.StateMachineInterceptorAdapter; +import org.springframework.stereotype.Component; +import ru.clearing.classes.statics.data.misc.Session; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.session.stage.SecondaryAuctionT0Session; +import ru.spcex.clearing.session.stage.TaskType; +import ru.spcex.clearing.session.state.DataEnum; +import ru.spcex.clearing.session.state.SsnEvent; +import ru.spcex.clearing.util.StateMachineUtil; +import ru.spcex.platform.enumeration.SessionStatus; +import ru.spcex.platform.enumeration.WorkflowStatus; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; +import ru.spcex.platform.utils.log.ExceptionUtils; + +/** + * при каждом переходе в другое состояние меняет статус сессии + */ +@Component +public class SessionStatusChangingInterceptor extends StateMachineInterceptorAdapter { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final Imdg ssnImdg; + + @Autowired + public SessionStatusChangingInterceptor(ImdgProvider imdgProvider, SecondaryAuctionT0Session secondaryAuctionT0Session) { + this.ssnImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Session, Session.class); + } + + + @Override + public StateContext preTransition(StateContext ctx) { + String id = StateMachineUtil.getId(ctx); + try { + State source = ctx.getTarget(); + State target = ctx.getTarget(); + if (target == null) { + throw excp( + ctx, + id + " transition target cannot be empty" + ); + } + log.info("{} preTransition interceptor: source={}/target={}", + id, + source != null ? source.getId() : "", + target.getId()); + + Session session = ctx.getExtendedState().get(DataEnum.session, Session.class); + + if (session == null && target.getId() == TaskType.StartRevise) { + return ctx; + } else if (session == null) { + throw excp(ctx, "%s source=%s/target=%s and session is not inside the extended state" + .formatted( + id, + source != null ? source.getId() : "unknown", + target.getId() + )); + } + + String sessionStatus = target.getId().getKey(); + session.setSessionStatus(sessionStatus); + ssnImdg.update(session); + return ctx; + } catch (Exception e) { + if (!ctx.getStateMachine().hasStateMachineError()) { + ctx.getStateMachine().setStateMachineError(e); + } + log.error("{} {}", id, ExceptionUtils.getStackTrace(e)); + return null; + } + } + + @Override + public Exception stateMachineError( + StateMachine stM, + Exception exception + ) { + String id = StateMachineUtil.getId(stM); + log.error("{} interceptor caught an exception!!!", id); + Session session = stM.getExtendedState().get(DataEnum.session, Session.class); + if (session != null) { + log.info("{} updating session status to error", id); + session.setSessionStatus(SessionStatus.CLOS.getKey()); + session.setWorkflowStatus(WorkflowStatus.Blocked.getKey()); + session.setUpdated(Instant.now()); + ssnImdg.update(session); + } else { + log.warn("{} no session in context", id); + } + + return exception; + } + + private IllegalStateException excp(StateContext ctx, String msg) { + IllegalStateException excp = new IllegalStateException(msg); + ctx.getStateMachine().setStateMachineError(excp); + return excp; + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/listener/MachineStopListener.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/listener/MachineStopListener.java new file mode 100644 index 000000000..d52b6d6b1 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/listener/MachineStopListener.java @@ -0,0 +1,43 @@ +package ru.spcex.clearing.session.state.listener; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.statemachine.StateMachine; +import org.springframework.statemachine.listener.StateMachineListenerAdapter; +import org.springframework.statemachine.state.PseudoStateKind; +import org.springframework.statemachine.state.State; +import ru.spcex.clearing.session.stage.TaskType; +import ru.spcex.clearing.session.state.SsnEvent; + +public class MachineStopListener extends StateMachineListenerAdapter { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final Runnable clearCurrentSession; + + public MachineStopListener(Runnable clearCurrentSession) { + this.clearCurrentSession = clearCurrentSession; + } + + /** + * видимо лишнее, т.к. end states и так останавливают машину + */ + @Override + public void stateChanged(State from, State to) { + if (to != null + && to.getId() != null + && to.getPseudoState() != null + && to.getPseudoState().getKind() != null + ) { + PseudoStateKind kind = to.getPseudoState().getKind(); + if (kind.equals(PseudoStateKind.END)) { + clearCurrentSession.run(); + } + } + } + + + @Override + public void stateMachineError(StateMachine stateMachine, Exception exception) { + log.error("SM LISTENER state machine encountered an error: {}", exception.getMessage()); + clearCurrentSession.run(); + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/util/StateMachineUtil.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/util/StateMachineUtil.java new file mode 100644 index 000000000..0a652d293 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/util/StateMachineUtil.java @@ -0,0 +1,28 @@ +package ru.spcex.clearing.util; + +import org.springframework.statemachine.StateContext; +import org.springframework.statemachine.StateMachine; +import org.springframework.statemachine.action.Action; + +public class StateMachineUtil { + + public static String getId(StateMachine machine) { + String id = machine.getId(); + if (id != null) return id; + return "[unknown SM]"; + } + + public static String getId(StateContext machine) { + StateMachine stateMachine = machine.getStateMachine(); + if (stateMachine == null) return "[unknown SM]"; + return getId(machine); + } + + public static Action chain(Action a, Action b) { + return context -> { + a.execute(context); + b.execute(context); + }; + } + +} diff --git a/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/StateMachineTest.java b/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/StateMachineTest.java new file mode 100644 index 000000000..85b32b3ce --- /dev/null +++ b/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/StateMachineTest.java @@ -0,0 +1,111 @@ +package ru.spcex.clearing.session.teststate; + +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.messaging.Message; +import org.springframework.statemachine.StateContext; +import org.springframework.statemachine.StateMachine; +import org.springframework.statemachine.config.StateMachineFactory; +import org.springframework.statemachine.support.StateMachineInterceptorAdapter; +import org.springframework.statemachine.transition.Transition; +import ru.spcex.clearing.session.teststate.config.Event; +import ru.spcex.clearing.session.teststate.config.State; +import ru.spcex.clearing.session.teststate.config.TestStateMachineConfig; +import ru.spcex.clearing.session.teststate.config.TestStateMachineExecutorsConfig; + +@SpringBootTest(classes = { + TestStateMachineConfig.class, + TestStateMachineExecutorsConfig.class +}) +class StateMachineTest { + private final Logger log = LoggerFactory.getLogger(getClass()); + + @Autowired + private StateMachineFactory stateMachineFactory; + + private StateMachine stateMachine; + + + @BeforeEach + public void setup() { + // Create a new state machine for each test + stateMachine = stateMachineFactory.getStateMachine(); + stateMachine.start(); + } + + @Test + void test() { + log.debug("test???"); + stateMachine + .getStateMachineAccessor() + .doWithRegion(function -> function.addStateMachineInterceptor( + new StateMachineInterceptorAdapter<>() { + @Override + public Message preEvent(Message message, StateMachine stateMachine) { + Event payload = message.getPayload(); + log.info("INTERCEPTOR preEvent catched event {}", payload); + return super.preEvent(message, stateMachine); + } + + @Override + public void preStateChange(org.springframework.statemachine.state.State state, Message message, Transition transition, StateMachine stateMachine, StateMachine rootStateMachine) { + log.info("INTERCEPTOR preStateChange catched"); + super.preStateChange(state, message, transition, stateMachine, rootStateMachine); + } + + @Override + public void postStateChange(org.springframework.statemachine.state.State state, Message message, Transition transition, StateMachine stateMachine, StateMachine rootStateMachine) { + log.info("INTERCEPTOR postStateChange catched"); + super.postStateChange(state, message, transition, stateMachine, rootStateMachine); + } + + @Override + public StateContext preTransition(StateContext stateContext) { + log.info("INTERCEPTOR preTransition catched"); + org.springframework.statemachine.state.State target = stateContext.getTarget(); + //if (target.getId().equals(State.continueRevise)) { + // IllegalStateException ex = new IllegalStateException("pre transition interceptor exception!!!"); + // stateContext.getStateMachine().setStateMachineError(ex); + // throw ex; + //} + return super.preTransition(stateContext); + } + + @Override + public StateContext postTransition(StateContext stateContext) { + log.info("INTERCEPTOR postTransition catched"); + return super.postTransition(stateContext); + } + + @Override + public Exception stateMachineError(StateMachine stateMachine, Exception exception) { + log.info("INTERCEPTOR error caught!!!"); + return exception; + } + })); + + stateMachine.sendEvent(Event.startSession); + try { + Thread.sleep(10000); + } catch (InterruptedException e) { + throw new RuntimeException(e); + } +// stateMachine.sendEvent(Event.continueRevise); +// log.info("'ve send a continueRevise"); +// try { +// Thread.sleep(100010); +// } catch (InterruptedException e) { +// throw new RuntimeException(e); +// } +// log.info("FIRST STOP"); +// stateMachine.stop(); +// log.info("SECOND STOP"); +// stateMachine.stop(); + log.info("So. We end with SM in state: {}", stateMachine.getState().getId());; + } + +} \ No newline at end of file diff --git a/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/Event.java b/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/Event.java new file mode 100644 index 000000000..6fe24c386 --- /dev/null +++ b/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/Event.java @@ -0,0 +1,5 @@ +package ru.spcex.clearing.session.teststate.config; + +public enum Event { + startSession, continueRevise, +} diff --git a/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/State.java b/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/State.java new file mode 100644 index 000000000..632fb4337 --- /dev/null +++ b/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/State.java @@ -0,0 +1,8 @@ +package ru.spcex.clearing.session.teststate.config; + +public enum State { + initial, + startRevise, + continueRevise, + newState +} diff --git a/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/TestStateMachineConfig.java b/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/TestStateMachineConfig.java new file mode 100644 index 000000000..c0c993427 --- /dev/null +++ b/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/TestStateMachineConfig.java @@ -0,0 +1,150 @@ +package ru.spcex.clearing.session.teststate.config; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.context.annotation.Configuration; +import org.springframework.core.task.TaskExecutor; +import org.springframework.messaging.Message; +import org.springframework.statemachine.StateMachine; +import org.springframework.statemachine.config.EnableStateMachineFactory; +import org.springframework.statemachine.config.EnumStateMachineConfigurerAdapter; +import org.springframework.statemachine.config.builders.StateMachineConfigurationConfigurer; +import org.springframework.statemachine.config.builders.StateMachineStateConfigurer; +import org.springframework.statemachine.config.builders.StateMachineTransitionConfigurer; +import org.springframework.statemachine.listener.StateMachineListenerAdapter; +import org.springframework.statemachine.transition.Transition; +import ru.spcex.clearing.session.teststate.config.action.ContinueReviseTestAction; +import ru.spcex.clearing.session.teststate.config.action.InitialAction; +import ru.spcex.clearing.session.teststate.config.action.SelfAction; +import ru.spcex.clearing.session.teststate.config.action.Send56TestAction; +import ru.spcex.clearing.session.teststate.config.guard.NoActiveSessionGuard; +import ru.spcex.clearing.util.StateMachineUtil; +import ru.spcex.platform.utils.log.ExceptionUtils; + +@Configuration +@EnableStateMachineFactory +public class TestStateMachineConfig extends EnumStateMachineConfigurerAdapter { + private final Logger log = LoggerFactory.getLogger(getClass()); + + @Autowired + @Qualifier("myStateMachineTaskExecutor") + private TaskExecutor taskExecutor; + + + @Override + public void configure(StateMachineStateConfigurer states) throws Exception { + states + .withStates() + .initial(State.initial, StateMachineUtil.chain(new InitialAction(), new SelfAction())) + .state(State.initial) + .state(State.startRevise) + .state(State.continueRevise) + .state(State.newState) + ; + + } + + @Override + public void configure(StateMachineTransitionConfigurer transitions) throws Exception { + transitions + .withExternal() + .source(State.initial) + .target(State.startRevise) + .event(Event.startSession) + .guard(new NoActiveSessionGuard()) + .action(new Send56TestAction()) + .and() + .withExternal() + .source(State.startRevise) + .target(State.continueRevise) + .action(new ContinueReviseTestAction()) + + //.and() + //.withExternal() + // .source(State.continueRevise) + // .target(State.newState) + // .action(context -> log.info("action without an event")) + .and() + .withExternal() + .source(State.continueRevise) + .target(State.newState) + .action(context -> { + log.info("action without an event"); + }) + .and() + .withInternal() + .source(State.newState) + .action(context -> { + log.info("INTERNAL ACTION WITHOUT AN EVENT"); + }) + ; + } + + @Override + public void configure(StateMachineConfigurationConfigurer config) throws Exception { + StateMachineListenerAdapter loggingChangeStateListener = new StateMachineListenerAdapter<>() { + @Override + public void stateEntered(org.springframework.statemachine.state.State state) { + State enteredState = state != null ? state.getId() : null; + log.info(String.format("LISTENER stateEntered: %s", enteredState)); + } + + @Override + public void eventNotAccepted(Message event) { + Event payload = event != null ? event.getPayload() : null; + log.info(String.format("LISTENER eventNotAccepted: %s", payload)); + } + + @Override + public void transition(Transition transition) { + State source = transition.getSource() != null ? transition.getSource().getId() : null; + State target = transition.getTarget().getId(); + log.info("LISTENER transition: source {} target {}", + source, target); + } + + @Override + public void stateChanged(org.springframework.statemachine.state.State from, org.springframework.statemachine.state.State to) { + State source = from != null ? from.getId() : null; + State target = to.getId(); + log.info("LISTENER stateChanged: source {} target {}", + source, target); + } + + @Override + public void stateExited(org.springframework.statemachine.state.State state) { + State whichOne = state != null ? state.getId() : null; + log.info("LISTENER stateExited: {}", whichOne); + } + + @Override + public void transitionEnded(Transition transition) { + State source = transition.getSource() != null ? transition.getSource().getId() : null; + State target = transition.getTarget().getId(); + log.info("LISTENER transitionEnded: source {} target {}", + source, target); + } + + @Override + public void stateMachineError(StateMachine stateMachine, Exception exception) { + log.info("LISTENER stateMachineError: {}", (ExceptionUtils.getStackTrace(exception))); + } + + @Override + public void transitionStarted(Transition transition) { + State source = transition.getSource() != null ? transition.getSource().getId() : null; + State target = transition.getTarget().getId(); + log.info("LISTENER transitionStarted: source {} target {}", + source, target); + } + }; + config + .withConfiguration() + .listener(loggingChangeStateListener) +// .taskExecutor(taskExecutor) + ; + + } +} diff --git a/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/TestStateMachineExecutorsConfig.java b/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/TestStateMachineExecutorsConfig.java new file mode 100644 index 000000000..3064bb251 --- /dev/null +++ b/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/TestStateMachineExecutorsConfig.java @@ -0,0 +1,32 @@ +package ru.spcex.clearing.session.teststate.config; + +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.core.task.TaskExecutor; +import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; + +@Configuration +public class TestStateMachineExecutorsConfig { + + @Bean(name = "myStateMachineTaskExecutor") + public TaskExecutor taskExecutor() { + ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); + executor.setCorePoolSize(1); + executor.setMaxPoolSize(1); + executor.setQueueCapacity(100); + executor.setThreadNamePrefix("state-machine-thread-"); + executor.initialize(); + return executor; + } + +// @Autowired +// @Qualifier("myStateMachineTaskScheduler") +// private TaskScheduler taskScheduler; +// +// @Bean(name = "myStateMachineTaskScheduler") +// public TaskScheduler taskScheduler() { +// final ThreadPoolTaskScheduler scheduler = new ThreadPoolTaskScheduler(); +// scheduler.setPoolSize(1); +// return scheduler; +// } +} diff --git a/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/action/ContinueReviseTestAction.java b/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/action/ContinueReviseTestAction.java new file mode 100644 index 000000000..61308a5d6 --- /dev/null +++ b/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/action/ContinueReviseTestAction.java @@ -0,0 +1,21 @@ +package ru.spcex.clearing.session.teststate.config.action; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.statemachine.StateContext; +import org.springframework.statemachine.action.Action; +import ru.spcex.clearing.session.teststate.config.Event; +import ru.spcex.clearing.session.teststate.config.State; + +public class ContinueReviseTestAction implements Action { + private final Logger log = LoggerFactory.getLogger(getClass()); + + @Override + public void execute(StateContext context) { + log.info("ACTION ContinueReviseTestAction"); +// IllegalStateException internalStageError = new IllegalStateException("internal stage error"); +// context.getStateMachine().setStateMachineError(internalStageError); + //throw internalStageError; + } + +} diff --git a/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/action/InitialAction.java b/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/action/InitialAction.java new file mode 100644 index 000000000..811a31b73 --- /dev/null +++ b/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/action/InitialAction.java @@ -0,0 +1,21 @@ +package ru.spcex.clearing.session.teststate.config.action; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.statemachine.StateContext; +import org.springframework.statemachine.action.Action; +import ru.spcex.clearing.session.teststate.config.Event; +import ru.spcex.clearing.session.teststate.config.State; + +public class InitialAction implements Action { + private final Logger log = LoggerFactory.getLogger(getClass()); + + @Override + public void execute(StateContext context) { + log.info("ACTION InitialAction"); +// IllegalStateException internalStageError = new IllegalStateException("internal stage error"); +// context.getStateMachine().setStateMachineError(internalStageError); + //throw internalStageError; + } + +} diff --git a/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/action/SelfAction.java b/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/action/SelfAction.java new file mode 100644 index 000000000..22653b70b --- /dev/null +++ b/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/action/SelfAction.java @@ -0,0 +1,21 @@ +package ru.spcex.clearing.session.teststate.config.action; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.statemachine.StateContext; +import org.springframework.statemachine.action.Action; +import ru.spcex.clearing.session.teststate.config.Event; +import ru.spcex.clearing.session.teststate.config.State; + +public class SelfAction implements Action { + private final Logger log = LoggerFactory.getLogger(getClass()); + + @Override + public void execute(StateContext context) { + log.info("ACTION SelfAction."); +// IllegalStateException internalStageError = new IllegalStateException("internal stage error"); +// context.getStateMachine().setStateMachineError(internalStageError); + //throw internalStageError; + } + +} diff --git a/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/action/Send56TestAction.java b/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/action/Send56TestAction.java new file mode 100644 index 000000000..1803c6f3f --- /dev/null +++ b/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/action/Send56TestAction.java @@ -0,0 +1,16 @@ +package ru.spcex.clearing.session.teststate.config.action; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.statemachine.StateContext; +import org.springframework.statemachine.action.Action; +import ru.spcex.clearing.session.teststate.config.Event; +import ru.spcex.clearing.session.teststate.config.State; + +public class Send56TestAction implements Action { + private final Logger log = LoggerFactory.getLogger(getClass()); + @Override + public void execute(StateContext context) { + log.info("ACTION Send56TestAction"); + } +} diff --git a/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/guard/NoActiveSessionGuard.java b/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/guard/NoActiveSessionGuard.java new file mode 100644 index 000000000..317aa5b0d --- /dev/null +++ b/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/session/teststate/config/guard/NoActiveSessionGuard.java @@ -0,0 +1,17 @@ +package ru.spcex.clearing.session.teststate.config.guard; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.statemachine.StateContext; +import org.springframework.statemachine.guard.Guard; +import ru.spcex.clearing.session.teststate.config.Event; +import ru.spcex.clearing.session.teststate.config.State; + +public class NoActiveSessionGuard implements Guard { + private final Logger log = LoggerFactory.getLogger(getClass()); + @Override + public boolean evaluate(StateContext context) { + log.info("***** NoActiveSessionGuard started"); + return true; + } +} diff --git a/clearing-parent/clearing-service/src/test/resources/logback-test.xml b/clearing-parent/clearing-service/src/test/resources/logback-test.xml new file mode 100644 index 000000000..c8732c152 --- /dev/null +++ b/clearing-parent/clearing-service/src/test/resources/logback-test.xml @@ -0,0 +1,19 @@ + + + + + + + ${CONSOLE_LOG_PATTERN} + utf-8 + + + + + + + + + + + \ No newline at end of file