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