Revert "пример spring state machine для FinalSession, тест, понизил версию"

This reverts commit 2e60049163.
This commit is contained in:
ialbert 2024-09-23 19:08:43 +03:00
parent 2e60049163
commit 4c804d243e
35 changed files with 23 additions and 2840 deletions

View file

@ -20,7 +20,7 @@
<dependency> <dependency>
<groupId>org.springframework.statemachine</groupId> <groupId>org.springframework.statemachine</groupId>
<artifactId>spring-statemachine-starter</artifactId> <artifactId>spring-statemachine-starter</artifactId>
<version>2.5.1</version> <version>3.2.0</version>
</dependency> </dependency>
<dependency> <dependency>
<groupId>ru.spcex.clearing</groupId> <groupId>ru.spcex.clearing</groupId>

View file

@ -1,11 +1,9 @@
package ru.spcex.clearing.config.state_machine_1; package ru.spcex.clearing.config.session;
import org.slf4j.Logger; import org.slf4j.Logger;
import org.slf4j.LoggerFactory; import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.InitializingBean; import org.springframework.beans.factory.InitializingBean;
import org.springframework.beans.factory.annotation.Autowired; 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.StateMachine;
import org.springframework.statemachine.config.StateMachineFactory; import org.springframework.statemachine.config.StateMachineFactory;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
@ -25,16 +23,11 @@ public class MachineTestService implements InitializingBean {
@Override @Override
public void afterPropertiesSet() throws Exception { public void afterPropertiesSet() throws Exception {
stateMachine.start(); // stateMachine.start();
log.info("send start revise event"); // log.info("send start revise event");
stateMachine.sendEvent(SessionEvent.Revise); // stateMachine.sendEvent(SessionEvent.Revise);
log.info("send continue revise event"); // log.info("send continue revise event");
stateMachine.sendEvent(SessionEvent.SdfReceived); // stateMachine.sendEvent(SessionEvent.SdfReceived);
Message<SessionEvent> eventMessage = MessageBuilder // log.info("end test");
.withPayload(SessionEvent.SdfReceived)
.setHeader("myPayload", new Object()) // Set the payload in the message header
.build();
stateMachine.sendEvent(eventMessage);
log.info("end test");
} }
} }

View file

@ -1,4 +1,4 @@
package ru.spcex.clearing.config.state_machine_1; package ru.spcex.clearing.config.session;
public enum SessionEvent { public enum SessionEvent {
Revise, Revise,

View file

@ -1,4 +1,4 @@
package ru.spcex.clearing.config.state_machine_1; package ru.spcex.clearing.config.session;
public enum StageDataEnum { public enum StageDataEnum {
sessionId, session, ExecutionList, PaymentInstructions sessionId, session, ExecutionList, PaymentInstructions

View file

@ -1,6 +1,5 @@
package ru.spcex.clearing.config.state_machine_1; package ru.spcex.clearing.config.session;
import java.util.function.Function;
import org.springframework.statemachine.ExtendedState; import org.springframework.statemachine.ExtendedState;
import org.springframework.statemachine.StateContext; import org.springframework.statemachine.StateContext;
import org.springframework.statemachine.action.Action; import org.springframework.statemachine.action.Action;
@ -9,6 +8,8 @@ import ru.spcex.clearing.session.stage.StageResult;
import ru.spcex.clearing.session.stage.Task; import ru.spcex.clearing.session.stage.Task;
import ru.spcex.clearing.session.stage.TaskType; import ru.spcex.clearing.session.stage.TaskType;
import java.util.function.Function;
/** /**
* адаптеры между ISessionStage старого образца к Spring State Machine<br> * адаптеры между ISessionStage старого образца к Spring State Machine<br>
* главное назначение: это сохранение результатов работы старого ISessionStage в ExtendedState<br> * главное назначение: это сохранение результатов работы старого ISessionStage в ExtendedState<br>
@ -19,7 +20,7 @@ public class StateActionAdapter implements Action<TaskType, SessionEvent> {
private Function<ExtendedState, Object> payloadForStageGetter; private Function<ExtendedState, Object> payloadForStageGetter;
private StageDataEnum saveName; private StageDataEnum saveName;
public StateActionAdapter(ISessionStage session) { protected StateActionAdapter(ISessionStage session) {
this.session = session; this.session = session;
} }

View file

@ -1,11 +1,5 @@
package ru.spcex.clearing.config.state_machine_1; package ru.spcex.clearing.config.session;
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.Logger;
import org.slf4j.LoggerFactory; import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.beans.factory.annotation.Qualifier;
@ -21,17 +15,7 @@ import ru.clearing.classes.statics.data.execution.ExecutionFond;
import ru.clearing.classes.statics.data.misc.Session; import ru.clearing.classes.statics.data.misc.Session;
import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.session.stage.TaskType; import ru.spcex.clearing.session.stage.TaskType;
import ru.spcex.clearing.session.stage.impl.BalanceRevise; import ru.spcex.clearing.session.stage.impl.*;
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.DealsPreparePayload;
import ru.spcex.clearing.session.stage.task.FormingPaymentInstructionPayload; import ru.spcex.clearing.session.stage.task.FormingPaymentInstructionPayload;
import ru.spcex.clearing.session.stage.task.InclusionToPoolPayload; import ru.spcex.clearing.session.stage.task.InclusionToPoolPayload;
@ -44,6 +28,13 @@ import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider; import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder; 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") @Configuration("sessionStateMachineFactory")
@EnableStateMachineFactory(name = "PrimaryBnSessionStateMachineFactory") @EnableStateMachineFactory(name = "PrimaryBnSessionStateMachineFactory")
public class StateBnConfig extends EnumStateMachineConfigurerAdapter<TaskType, SessionEvent> { public class StateBnConfig extends EnumStateMachineConfigurerAdapter<TaskType, SessionEvent> {

View file

@ -1,173 +0,0 @@
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<TaskType, SsnEvent> {
private final Logger log = LoggerFactory.getLogger(getClass());
private final Imdg<Session> sessionImdg;
private final Imdg<Registry> 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<List<String>> 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<TaskType, SsnEvent> 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<TaskType, SsnEvent> 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<TaskType, SsnEvent> config) throws Exception {
StateMachineListenerAdapter<TaskType, SsnEvent> loggingChangeStateListener = new StateMachineListenerAdapter<>() {
@Override
public void stateEntered(State<TaskType, SsnEvent> state) {
TaskType enteredState = state != null ? state.getId() : null;
log.info("State entered: {}", enteredState);
}
};
config
.withConfiguration()
.machineId(SessionType.FINL.getKey())
.listener(loggingChangeStateListener)
;
}
}

View file

@ -1,5 +0,0 @@
package ru.spcex.clearing.session.state;
public enum ActionHeader {
sdf56FromTime;
}

View file

@ -1,9 +0,0 @@
package ru.spcex.clearing.session.state;
public enum DataEnum {
sessionId, //Long
sessionType, //Session
session, //Session
dealsPrepared, //List<ExecutionCommon>
counterPartyId, //Long
}

View file

@ -1,7 +0,0 @@
package ru.spcex.clearing.session.state;
public enum SsnEvent {
Revise,
Sdf57Processed,
SdfReceived,
}

View file

@ -1,36 +0,0 @@
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<TaskType, SsnEvent> {
private final Logger log = LoggerFactory.getLogger(getClass());
@Override
public void execute(StateContext<TaskType, SsnEvent> 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<TaskType, SsnEvent> context);
}

View file

@ -1,43 +0,0 @@
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<Session> sessionImdg;
private final Logger log = LoggerFactory.getLogger(CreateSessionAction.class);
public CreateSessionAction(Imdg<Session> sessionImdg) {
this.sessionImdg = sessionImdg;
}
@Override
public void actualExecute(StateContext<TaskType, SsnEvent> 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);
}
}

View file

@ -1,140 +0,0 @@
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<ExecutionDeposit> executionDepositImdg;
private final Imdg<ExecutionFond> executionFondImdg;
private final Imdg<ExecutionCurrency> executionCurrencyImdg;
private final List<ImdgPredicate> execDepositPredicates = new ArrayList<>();
private final List<ImdgPredicate> execFondPredicates = new ArrayList<>();
private final List<ImdgPredicate> 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<TaskType, SsnEvent> 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<List<ImdgPredicate>, 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<ExecutionCommon> 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 <E extends ExecutionCommon> Function<E, ExecutionCommon> execToInterface() {
return (e) -> e;
}
}

View file

@ -1,245 +0,0 @@
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<Registry> registryImdg;
private Imdg<Session> sessionImdg;
private final Imdg<ExecutionDeposit> executionDepositImdg;
private final Imdg<ExecutionFond> executionFondImdg;
private final Imdg<ExecutionCurrency> executionCurrImdg;
private KafkaSender kafkaSender;
private final List<ImdgPredicate> 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<TaskType, SsnEvent> 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<Registry> obligations = registryImdg.getCollectionObjectsByPredicate(registryPredicate);
Map<Long, List<Registry>> registryByGroupId = obligations.stream()
.filter(registry -> registry.getSettlementDate().isEqual(LocalDate.now()))
.collect(Collectors.groupingBy(Registry::getGroupId));
Optional<SessionType> ssnType = Optional.of(sessionType);
Map<Long, Registry> rgsToUpdate = new HashMap<>();
List<ExecutionCommon> execsToUpdate = new ArrayList<>();
boolean loadExecs = !IEnumKey.contains(sessionType,
SessionType.TRDT, SessionType.CURR, SessionType.UNIT, SessionType.IPOT);
for (Map.Entry<Long, List<Registry>> 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 <T extends ExecutionCommon> Collection<ExecutionCommon> 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<T> execImdg = null;
switch (ssnTpe) {
case MEDM, FINL, XDEP -> execImdg = (Imdg<T>) executionDepositImdg;
case IPO0, IPOB, TRDT, IPOT -> execImdg = (Imdg<T>) executionFondImdg;
case CURR -> execImdg = (Imdg<T>) executionCurrImdg;
default -> {
//кейс для "общих" сессий (напр. UNIT) в которых сочетаются разные сделки
//смотрим на секцию регистра.
Section section;
if ((section = IEnumKey.getEnumByKey(Section.class, rgsSection)) != null) {
if (section.equals(Section.MKR)) execImdg = (Imdg<T>) executionDepositImdg;
else if (section.equals(Section.FOND)) execImdg = (Imdg<T>) executionFondImdg;
else if (section.equals(Section.CURR)) execImdg = (Imdg<T>) 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<T> execs = execImdg.getCollectionObjectsByPredicate(prdct);
//Imdg<T> finalExecImdg = execImdg;
execs.forEach(e -> {
e.setUpdated(now);
e.setSessionId(sessionId);
// finalExecImdg.update(e);
});
return (Collection<ExecutionCommon>) 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<ExecutionCommon> execs) {
log.debug("updating {} executions", execs.size());
execs.sort(Comparator.comparing(ExecutionCommon::type));
log.debug("sorted executions by type");
Map<Long, ExecutionCommon> 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();
}
}
}

View file

@ -1,285 +0,0 @@
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<Registry> registryImdg;
private final Imdg<ClearingMemberCategory> 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<TaskType, SsnEvent> 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<Map.Entry<Long, List<Registry>>> 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<Long, List<Registry>> entry : registriesByGroupSorted) {
List<Registry> group = entry.getValue();
Optional<Registry> omtInGroupO = group.stream().filter(registry -> equalByRgs(OM_T, registry)).findFirst();
Optional<Registry> 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<Registry> 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<Registry> amfAssetO = searchAssetByOMT(omtRgs);
Optional<Registry> 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<Registry, BigDecimal> processAssetsAndSetSessionId = (amf, amount) -> {
Optional<AssetTrio> asts = assets.processByAm_f(amf, amount);
if (asts.isPresent()) {
asts.get().a__b().setSessionId(sessionId);
registryImdg.update(asts.get().a__b());
}
};
Optional<Registry> 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<Registry> 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<Registry> 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<Boolean> 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;
}
}

View file

@ -1,408 +0,0 @@
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<Registry> 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<TaskType, SsnEvent> 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<Registry> registriesToProcess = registryImdg.getCollectionObjectsBySQL(sqlCondition);
Map<Long, Registry> rgssToStore = new HashMap<>(registriesToProcess.size()); //вторые ноги у сделок - нужна ли оптимизация?
List<Map.Entry<Long, List<Registry>>> 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<Long, List<Registry>> entry : registriesByGroupSorted) {
List<Registry> group = entry.getValue();
List<Registry> 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<CheckResult> 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<Registry> 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<Boolean> 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<CheckResult> checkResults, List<Registry> registries, Map<Long, Registry> 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<Registry> 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<Registry> getRefundDateRgsIfPresent(List<Registry> registries) {
if (registries.stream().anyMatch(rgs -> rgs.getRefundDate() != null) && !SessionType.MEDM.equals(sessionType)) {
Optional<Pair<Long, LocalDate>> 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<Long, Registry> 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;
}
}
}

View file

@ -1,135 +0,0 @@
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<Registry> registryImdg;
private final Imdg<Relation> imdgRltn;
private final Imdg<Account> imdgAcc;
private final Imdg<Company> imdgCmp;
private final Imdg<TradingClearingRegistry> 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<TaskType, SsnEvent> ctx) {
try (cash) {
Long sessionId = ctx.getExtendedState().get(DataEnum.sessionId, Long.class);
obligationAdmission(sessionId);
}
}
private void obligationAdmission(Long sessionId) {
Collection<Registry> registries = registryImdg.getCollectionObjectsBySQL("sessionId = " + sessionId);
Map<Long, List<Registry>> byGroups = registries.stream().collect(Collectors.groupingBy(Registry::getGroupId));
log.info("found {} ({} groups) registries with sessionId {}", registries.size(), byGroups.size(), sessionId);
Map<Long, Registry> rgsToUpdate = new HashMap<>();
for (Map.Entry<Long, List<Registry>> grpEntry : byGroups.entrySet()) {
EnumMessage groupError = null;
List<Registry> rgsGroup = grpEntry.getValue();
for (Registry rgs : rgsGroup) {
IValidator validator = valFor(rgs);
Optional<EnumMessage> 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<Registry> 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<Relation> rltnCash;
CashV2ById<Account> accCash;
CashV2ById<Company> cmpCash;
CashV2ById<TradingClearingRegistry> 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"));
}
}
}

View file

@ -1,454 +0,0 @@
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<Registry> registryImdg;
private final ImdgId idGen;
private final ImdgProvider imdgProvider;
protected final RequirementsAndObligationCreationCash cash;
private final Imdg<Company> cmpImdg;
private final Imdg<TradingClearingRegistry> tcrImdg;
private final Imdg<Account> accImdg;
private final Imdg<Session> ssnImdg;
private final Imdg<ClearingAccount> clAccImdg;
private final Imdg<DepoAccount> dpAccImdg;
private final Imdg<CurrencyPairSecurity> currPairSecImdg;
private final Imdg<FixedIncomeSecurity> fixIncSecImdg;
private final Imdg<MoneyMarketSecurity> monMrktSecImdg;
private final Imdg<EquitySecurity> eqtySecImdg;
private final Imdg<CurrencyPairDictionary> currPairDictImdg;
private final Imdg<Currency> currImdg;
private final TcrSearcher tcrSearcher;
private final SecuritySelector<Security> 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<TradingClearingRegistryList> 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 <T extends Security> Imdg<Security> castSecImdg(Imdg<T> m) {
return (Imdg<Security>) m;
}
protected Boolean doRgsSearch = null;
@Override
public void actualExecute(StateContext<TaskType, SsnEvent> ctx) {
try (cash) {
@SuppressWarnings("unchecked")
List<ExecutionCommon> 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<ExecutionCommon> data, Long sessionId, Map<Long, Registry> rgsStorage) {
data.sort(tradeTimeComparator.thenComparing(execIdComparator));
Map<Long, Registry> newRgss = rgsStorage != null ? rgsStorage : storageMap(data);
//см. описание к #matchExecutions
for (int i = 0; i < data.size(); ) {
Pair<ExecutionCommon, ExecutionCommon> 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<BiFunction<ExecutionCommon, RegistryDesignation, Optional<Registry>>> regSearchMethod = new AtomicReference<>(
(exc, dsg) -> Optional.empty()
);
//обычный поиск
if (searchForRegisters(sessionId)) {
regSearchMethod.set(this::findReg);
}
TriConsumer<ExecutionCommon, RegistryDesignation, Boolean> 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<Long, Registry> storageMap(List<ExecutionCommon> 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<Registry> 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<Registry> 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<Registry> 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<ExecutionCommon, ExecutionCommon> matchExecutions(List<ExecutionCommon> 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<T1, T2, T3> {
void accept(T1 t1, T2 t2, T3 t3);
}
protected static class RequirementsAndObligationCreationCash extends CashCloser {
protected final CashV2<Long, Company> cmpCash;
protected final CashV2<Long, TradingClearingRegistry> tcrCash;
protected final CashV2<Long, Account> accCash;
protected final CashV2<Long, DepoAccount> dpAccCash;
protected final CashV2<String, Currency> currCash;
protected final CashV2<Long, ClearingAccount> clrAccCash;
protected final CashV2<Long, Session> ssnCash;
protected final CashV2<Long, CurrencyPairSecurity> currPairCash;
protected final CashV2<Long, FixedIncomeSecurity> fixIncSecCash;
protected final CashV2<Long, MoneyMarketSecurity> monMrktSecCash;
protected final CashV2<Long, EquitySecurity> eqtySecCash;
protected final CashV2<Long, CurrencyPairDictionary> currPairDictCash;
protected final CashV2ByIdAndString<TradingClearingRegistryList> 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())));
}
}
}

View file

@ -1,55 +0,0 @@
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<Registry> registryImdg;
@Autowired
public ReviseStage1(ImdgProvider imdgProvider) {
this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
}
@Override
public void actualExecute(StateContext<TaskType, SsnEvent> ctx) {
String sql = RegistryCodeSqlBuilder.getInstance(RegistryTradingParams.A__T).build();
Collection<Registry> regsAT = registryImdg.getCollectionObjectsBySQL(sql);
log.trace("Select {} registry's by query \"{}\" for revision step 1", regsAT.size(), sql);
Instant now = Instant.now();
int updateCount = 0;
for (Registry reg : regsAT) {
if (reg.getBalance() == null) {
log.debug("Registry[{}] with null balance", reg.getId());
} else {
if (!Objects.equals(reg.getPlanBalance(), reg.getBalance())) {
reg.setPlanBalance(reg.getBalance());
reg.setUpdated(now);
registryImdg.update(reg);
updateCount++;
}
}
}
log.debug("At revision stage 1 do updated {} of {} registers {}", updateCount, regsAT.size(), sql);
}
}

View file

@ -1,97 +0,0 @@
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<Statement> statementImdg;
private final Imdg<Currency> currencyImdg;
private final Imdg<SDf56> 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<TaskType, SsnEvent> 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<Long> 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());
}
}

View file

@ -1,110 +0,0 @@
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<SessionType, StateMachineFactory<TaskType, SsnEvent>> factories;
private final SessionStatusChangingInterceptor statusChangingInterceptor;
private StateMachine<TaskType, SsnEvent> currentSession;
@Autowired
public SessionStateMachineWrapper(
Map<String, StateMachineFactory<TaskType, SsnEvent>> 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<LauncherCommandRequest> 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<TaskType, SsnEvent> 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<SsnEvent> 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<TaskType, SsnEvent> getByType(SessionType sessionType) {
StateMachineFactory<TaskType, SsnEvent> factory = factories.get(sessionType);
if (factory == null) {
throw new NoSuchElementException("no state machine for session: " + sessionType);
}
return factory;
}
private StateMachine<TaskType, SsnEvent> build(
StateMachineFactory<TaskType, SsnEvent> stateMachineFactory
) {
StateMachine<TaskType, SsnEvent> stateMachine = stateMachineFactory.getStateMachine();
stateMachine
.getStateMachineAccessor()
.doWithAllRegions(
access -> access.addStateMachineInterceptor(statusChangingInterceptor)
);
return stateMachine;
}
}

View file

@ -1,108 +0,0 @@
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<TaskType, SsnEvent> {
private final Logger log = LoggerFactory.getLogger(getClass());
private final Imdg<Session> ssnImdg;
@Autowired
public SessionStatusChangingInterceptor(ImdgProvider imdgProvider, SecondaryAuctionT0Session secondaryAuctionT0Session) {
this.ssnImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Session, Session.class);
}
@Override
public StateContext<TaskType, SsnEvent> preTransition(StateContext<TaskType, SsnEvent> ctx) {
String id = StateMachineUtil.getId(ctx);
try {
State<TaskType, SsnEvent> source = ctx.getTarget();
State<TaskType, SsnEvent> 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<TaskType, SsnEvent> 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<TaskType, SsnEvent> ctx, String msg) {
IllegalStateException excp = new IllegalStateException(msg);
ctx.getStateMachine().setStateMachineError(excp);
return excp;
}
}

View file

@ -1,43 +0,0 @@
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<TaskType, SsnEvent> {
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<TaskType, SsnEvent> from, State<TaskType, SsnEvent> 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<TaskType, SsnEvent> stateMachine, Exception exception) {
log.error("SM LISTENER state machine encountered an error: {}", exception.getMessage());
clearCurrentSession.run();
}
}

View file

@ -1,28 +0,0 @@
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 <S, A> Action<S, A> chain(Action<S, A> a, Action<S, A> b) {
return context -> {
a.execute(context);
b.execute(context);
};
}
}

View file

@ -1,111 +0,0 @@
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<State, Event> stateMachineFactory;
private StateMachine<State, Event> 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<Event> preEvent(Message<Event> message, StateMachine<State, Event> 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, Event> state, Message<Event> message, Transition<State, Event> transition, StateMachine<State, Event> stateMachine, StateMachine<State, Event> rootStateMachine) {
log.info("INTERCEPTOR preStateChange catched");
super.preStateChange(state, message, transition, stateMachine, rootStateMachine);
}
@Override
public void postStateChange(org.springframework.statemachine.state.State<State, Event> state, Message<Event> message, Transition<State, Event> transition, StateMachine<State, Event> stateMachine, StateMachine<State, Event> rootStateMachine) {
log.info("INTERCEPTOR postStateChange catched");
super.postStateChange(state, message, transition, stateMachine, rootStateMachine);
}
@Override
public StateContext<State, Event> preTransition(StateContext<State, Event> stateContext) {
log.info("INTERCEPTOR preTransition catched");
org.springframework.statemachine.state.State<State, Event> 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<State, Event> postTransition(StateContext<State, Event> stateContext) {
log.info("INTERCEPTOR postTransition catched");
return super.postTransition(stateContext);
}
@Override
public Exception stateMachineError(StateMachine<State, Event> 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());;
}
}

View file

@ -1,5 +0,0 @@
package ru.spcex.clearing.session.teststate.config;
public enum Event {
startSession, continueRevise,
}

View file

@ -1,8 +0,0 @@
package ru.spcex.clearing.session.teststate.config;
public enum State {
initial,
startRevise,
continueRevise,
newState
}

View file

@ -1,150 +0,0 @@
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<State, Event> {
private final Logger log = LoggerFactory.getLogger(getClass());
@Autowired
@Qualifier("myStateMachineTaskExecutor")
private TaskExecutor taskExecutor;
@Override
public void configure(StateMachineStateConfigurer<State, Event> 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<State, Event> 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<State, Event> config) throws Exception {
StateMachineListenerAdapter<State, Event> loggingChangeStateListener = new StateMachineListenerAdapter<>() {
@Override
public void stateEntered(org.springframework.statemachine.state.State<State, Event> state) {
State enteredState = state != null ? state.getId() : null;
log.info(String.format("LISTENER stateEntered: %s", enteredState));
}
@Override
public void eventNotAccepted(Message<Event> event) {
Event payload = event != null ? event.getPayload() : null;
log.info(String.format("LISTENER eventNotAccepted: %s", payload));
}
@Override
public void transition(Transition<State, Event> 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<State, Event> from, org.springframework.statemachine.state.State<State, Event> 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, Event> state) {
State whichOne = state != null ? state.getId() : null;
log.info("LISTENER stateExited: {}", whichOne);
}
@Override
public void transitionEnded(Transition<State, Event> 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<State, Event> stateMachine, Exception exception) {
log.info("LISTENER stateMachineError: {}", (ExceptionUtils.getStackTrace(exception)));
}
@Override
public void transitionStarted(Transition<State, Event> 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)
;
}
}

View file

@ -1,32 +0,0 @@
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;
// }
}

View file

@ -1,21 +0,0 @@
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<State, Event> {
private final Logger log = LoggerFactory.getLogger(getClass());
@Override
public void execute(StateContext<State, Event> context) {
log.info("ACTION ContinueReviseTestAction");
// IllegalStateException internalStageError = new IllegalStateException("internal stage error");
// context.getStateMachine().setStateMachineError(internalStageError);
//throw internalStageError;
}
}

View file

@ -1,21 +0,0 @@
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<State, Event> {
private final Logger log = LoggerFactory.getLogger(getClass());
@Override
public void execute(StateContext<State, Event> context) {
log.info("ACTION InitialAction");
// IllegalStateException internalStageError = new IllegalStateException("internal stage error");
// context.getStateMachine().setStateMachineError(internalStageError);
//throw internalStageError;
}
}

View file

@ -1,21 +0,0 @@
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<State, Event> {
private final Logger log = LoggerFactory.getLogger(getClass());
@Override
public void execute(StateContext<State, Event> context) {
log.info("ACTION SelfAction.");
// IllegalStateException internalStageError = new IllegalStateException("internal stage error");
// context.getStateMachine().setStateMachineError(internalStageError);
//throw internalStageError;
}
}

View file

@ -1,16 +0,0 @@
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<State, Event> {
private final Logger log = LoggerFactory.getLogger(getClass());
@Override
public void execute(StateContext<State, Event> context) {
log.info("ACTION Send56TestAction");
}
}

View file

@ -1,17 +0,0 @@
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<State, Event> {
private final Logger log = LoggerFactory.getLogger(getClass());
@Override
public boolean evaluate(StateContext<State, Event> context) {
log.info("***** NoActiveSessionGuard started");
return true;
}
}

View file

@ -1,19 +0,0 @@
<?xml version="1.0" encoding="UTF-8"?>
<configuration>
<property name="CONSOLE_LOG_PATTERN" value="%date{HH:mm:ss.SSS} [%thread] %-5level %class{0}:%line - %message%n" />
<appender name="CONSOLE" class="ch.qos.logback.core.ConsoleAppender">
<encoder>
<!-- |%X{ru.nbch.scoring.web.logging.mdc_key}-->
<Pattern>${CONSOLE_LOG_PATTERN}</Pattern>
<charset>utf-8</charset>
</encoder>
</appender>
<root level="info">
<appender-ref ref="CONSOLE"/>
</root>
<logger name="ru.spcex" level="trace" additivity="false">
<appender-ref ref="CONSOLE"/>
</logger>
</configuration>