diff --git a/clearing-parent/clearing-service/pom.xml b/clearing-parent/clearing-service/pom.xml index fe421bfb5..96acb6e8d 100644 --- a/clearing-parent/clearing-service/pom.xml +++ b/clearing-parent/clearing-service/pom.xml @@ -17,6 +17,11 @@ + + org.springframework.statemachine + spring-statemachine-starter + 3.2.0 + ru.spcex.platform platform-messaging diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/session/MachineTestService.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/session/MachineTestService.java new file mode 100644 index 000000000..fe980bed3 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/session/MachineTestService.java @@ -0,0 +1,33 @@ +package ru.spcex.clearing.config.session; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.InitializingBean; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.statemachine.StateMachine; +import org.springframework.statemachine.config.StateMachineFactory; +import org.springframework.stereotype.Component; +import ru.spcex.clearing.session.stage.TaskType; + +@Component +public class MachineTestService implements InitializingBean { + private final Logger log = LoggerFactory.getLogger(getClass()); + + private final StateMachine stateMachine; + + @Autowired + public MachineTestService(StateMachineFactory stateMachineFactory) { + this.stateMachine = stateMachineFactory.getStateMachine(); + } + + + @Override + public void afterPropertiesSet() throws Exception { +// stateMachine.start(); +// log.info("send start revise event"); +// stateMachine.sendEvent(SessionEvent.Revise); +// log.info("send continue revise event"); +// stateMachine.sendEvent(SessionEvent.SdfReceived); +// log.info("end test"); + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/session/SessionEvent.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/session/SessionEvent.java new file mode 100644 index 000000000..50430df1a --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/session/SessionEvent.java @@ -0,0 +1,6 @@ +package ru.spcex.clearing.config.session; + +public enum SessionEvent { + Revise, + SdfReceived, +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/session/StageDataEnum.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/session/StageDataEnum.java new file mode 100644 index 000000000..c12287a21 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/session/StageDataEnum.java @@ -0,0 +1,5 @@ +package ru.spcex.clearing.config.session; + +public enum StageDataEnum { + sessionId, session, ExecutionList, PaymentInstructions +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/session/StateActionAdapter.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/session/StateActionAdapter.java new file mode 100644 index 000000000..e0099ecfd --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/session/StateActionAdapter.java @@ -0,0 +1,50 @@ +package ru.spcex.clearing.config.session; + +import org.springframework.statemachine.ExtendedState; +import org.springframework.statemachine.StateContext; +import org.springframework.statemachine.action.Action; +import ru.spcex.clearing.session.stage.ISessionStage; +import ru.spcex.clearing.session.stage.StageResult; +import ru.spcex.clearing.session.stage.Task; +import ru.spcex.clearing.session.stage.TaskType; + +import java.util.function.Function; + +/** + * адаптеры между ISessionStage старого образца к Spring State Machine
+ * главное назначение: это сохранение результатов работы старого ISessionStage в ExtendedState
+ * + */ +public class StateActionAdapter implements Action { + private final ISessionStage session; + private Function payloadForStageGetter; + private StageDataEnum saveName; + + protected StateActionAdapter(ISessionStage session) { + this.session = session; + } + + public void setPayloadForStageGetter(Function payloadForStageGetter) { + this.payloadForStageGetter = payloadForStageGetter; + } + + public void setSaveName(StageDataEnum saveName) { + this.saveName = saveName; + } + + @Override + public void execute(StateContext context) { + TaskType currState = context.getStateMachine().getState().getId(); + Task task; + if (payloadForStageGetter != null) { + task = new Task<>(currState, payloadForStageGetter.apply(context.getExtendedState())); + } else { + task = new Task<>(currState, null); + } + StageResult submit = session.submit(task); + Object stageResult = submit.getStageResult(); + if (stageResult != null) { + context.getExtendedState().getVariables().put(saveName.name(), stageResult); + } + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/session/StateBnConfig.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/session/StateBnConfig.java new file mode 100644 index 000000000..2f337e091 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/session/StateBnConfig.java @@ -0,0 +1,218 @@ +package ru.spcex.clearing.config.session; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +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.execution.ExecutionFond; +import ru.clearing.classes.statics.data.misc.Session; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.session.stage.TaskType; +import ru.spcex.clearing.session.stage.impl.*; +import ru.spcex.clearing.session.stage.task.DealsPreparePayload; +import ru.spcex.clearing.session.stage.task.FormingPaymentInstructionPayload; +import ru.spcex.clearing.session.stage.task.InclusionToPoolPayload; +import ru.spcex.clearing.session.stage.task.InspectionPoolPayload; +import ru.spcex.platform.classes.base.interfaces.ExecutionType; +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; +import ru.spcex.platform.imdg.api.ImdgProvider; +import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder; + +import java.time.LocalDate; +import java.util.Arrays; +import java.util.HashSet; +import java.util.List; +import java.util.function.Supplier; + +@Configuration("sessionStateMachineFactory") +@EnableStateMachineFactory(name = "PrimaryBnSessionStateMachineFactory") +public class StateBnConfig extends EnumStateMachineConfigurerAdapter { + private Logger log = LoggerFactory.getLogger(getClass()); + + private final Imdg sessionImdg; + /** + * адаптеры между ISessionStage старого образца к Spring State Machine + */ + private final StateActionAdapter balanceReviseAction; + private final StateActionAdapter dealPrepareAction; + private final StateActionAdapter requirementsAndObligationCreationAction; + private final StateActionAdapter obligationAdmissionAction; + private final StateActionAdapter inclusionToPoolAction; + private final StateActionAdapter inspectionObligationsAction; + private final StateActionAdapter formingRegistersOnOSAction; + private final StateActionAdapter formingPaymentInstructionAction; + + public StateBnConfig(ImdgProvider imdgProvider, + BalanceRevise balanceRevise, + DealsPrepare dealsPrepare, + RequirementsAndObligationCreation requirementsAndObligationCreation, + ObligationAdmission obligationsAdmission, + InclusionObligations inclusionObligations, + InspectionObligations inspectionObligations, + FormingRegistersOnOS formingRegistersOnOS, + FormingPaymentInstruction formingPaymentInstruction, + UnlockResources unlockResources, + FinishingSession finishingSession, + EndStageNotification endStageNotification, + @Qualifier("marketCodesForBn") Supplier> marketCodes) { + + Imdg executionFondImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionFond, ExecutionFond.class); + this.sessionImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Session, Session.class); + + //настройка стадий + dealsPrepare.searchForExecutions(ExecutionType.ExecutionFond); + ImdgPredicateBuilder execFondPb = executionFondImdg.predicateBuilder(); + dealsPrepare.addExecutionFondCondition(execFondPb.regex("settlementCode", "^B\\d{2}$")); + dealsPrepare.addExecutionFondCondition(execFondPb.in("market", marketCodes.get().toArray(new String[0]))); + imdgProvider.waitAvailable(); + + //настройка адаптеров для State Machine + balanceReviseAction = new StateActionAdapter(balanceRevise); + + dealPrepareAction = new StateActionAdapter(dealsPrepare); + dealPrepareAction.setPayloadForStageGetter((extendedState) -> { + DealsPreparePayload payload = new DealsPreparePayload(); + payload.setSessionId(extendedState.get(StageDataEnum.sessionId.name(), Long.class)); + return payload; + }); + dealPrepareAction.setSaveName(StageDataEnum.ExecutionList); + + requirementsAndObligationCreationAction = new StateActionAdapter(requirementsAndObligationCreation); + requirementsAndObligationCreationAction.setPayloadForStageGetter((extendedState) + -> extendedState.get(StageDataEnum.ExecutionList.name(), List.class)); + + obligationAdmissionAction = new StateActionAdapter(obligationsAdmission); + obligationAdmissionAction.setPayloadForStageGetter((extendedState) -> extendedState.get(StageDataEnum.sessionId.name(), Long.class)); + + inclusionToPoolAction = new StateActionAdapter(inclusionObligations); + inclusionToPoolAction.setPayloadForStageGetter((extendedState) -> { + InclusionToPoolPayload payload = new InclusionToPoolPayload(); + Session session = extendedState.get(StageDataEnum.session.name(), Session.class); + payload.setSessionType(session.getSessionType()); + return payload; + }); + + inspectionObligationsAction = new StateActionAdapter(inspectionObligations); + inspectionObligationsAction.setPayloadForStageGetter((extendedState) -> { + InspectionPoolPayload payload = new InspectionPoolPayload(); + Session session = extendedState.get(StageDataEnum.session.name(), Session.class); + payload.setProcessedCompanyId(session.getCompanyId()); + return payload; + }); + + formingRegistersOnOSAction = new StateActionAdapter(formingRegistersOnOS); + + formingPaymentInstructionAction = new StateActionAdapter(formingPaymentInstruction); + formingPaymentInstructionAction.setPayloadForStageGetter((extendedState) -> { + FormingPaymentInstructionPayload payload = new FormingPaymentInstructionPayload(); + payload.setSessionId(extendedState.get(StageDataEnum.sessionId.name(), Long.class)); + return payload; + }); + formingPaymentInstructionAction.setSaveName(StageDataEnum.PaymentInstructions); + // if (paymentResult != null && paymentResult.getStageResult().isEmpty()) { + // runStage(TaskType.FormingPaymentInstruction, balanceRevise); + //// finishPart(req); + // } + } + + + @Override + public void configure(StateMachineStateConfigurer states) throws Exception { + states + .withStates() + .initial(TaskType.StartRevise, context -> { + Session newSession = new Session(); + newSession.setSection(Section.FOND.getKey()); + newSession.setSessionType(SessionType.IPOB.getKey()); + 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(StageDataEnum.sessionId.name(), newSession.getId()); + context.getExtendedState().getVariables().put(StageDataEnum.session.name(), newSession); + }) + .states(new HashSet<>(Arrays.asList(TaskType.StartRevise, + TaskType.ContinueRevise, + TaskType.DealsPrepare, + TaskType.RequirementsAndObligationsCreate, + TaskType.ObligationsAdmission, + TaskType.InclusionToPool, + TaskType.InspectionObligations, + TaskType.FormingRegistersOnOS, + TaskType.FormingPaymentInstruction +// TaskType.UnlockResources, +// TaskType.FinishingSession, +// TaskType.EndStageNotification + ))) + .end(TaskType.FormingPaymentInstruction) + ; + } + + + @Override + public void configure(StateMachineTransitionConfigurer transitions) throws Exception { + transitions.withExternal() + .event(SessionEvent.Revise) + .source(TaskType.StartRevise).target(TaskType.ContinueRevise) + .action(balanceReviseAction) + .and() + .withExternal() + .event(SessionEvent.SdfReceived) + .source(TaskType.ContinueRevise).target(TaskType.DealsPrepare) + .action(context -> { + log.info("SDF57 and SDF01 received, continue session"); + }) + .and() + .withExternal() + .source(TaskType.DealsPrepare).target(TaskType.RequirementsAndObligationsCreate) + .action(dealPrepareAction) + .and() + .withExternal() + .source(TaskType.RequirementsAndObligationsCreate).target(TaskType.ObligationsAdmission) + .action(requirementsAndObligationCreationAction) + .and() + .withExternal() + .source(TaskType.ObligationsAdmission).target(TaskType.InclusionToPool) + .action(obligationAdmissionAction) + .and() + .withExternal() + .source(TaskType.InclusionToPool).target(TaskType.InspectionObligations) + .action(inclusionToPoolAction) + .and() + .withExternal() + .source(TaskType.InspectionObligations).target(TaskType.FormingRegistersOnOS) + .action(inspectionObligationsAction) + .and() + .withExternal() + .source(TaskType.FormingRegistersOnOS).target(TaskType.FormingPaymentInstruction) + .action(formingRegistersOnOSAction) + .and() + .withExternal() + .source(TaskType.FormingPaymentInstruction).target(TaskType.FormingPaymentInstruction) + .action(formingPaymentInstructionAction); + + } + + @Override + public void configure(StateMachineConfigurationConfigurer config) throws Exception { + StateMachineListenerAdapter loggingChangeStateListener = new StateMachineListenerAdapter<>() { + @Override + public void stateEntered(State state) { + TaskType enteredState = state != null ? state.getId() : null; + log.info(String.format("State entered: %s", enteredState)); + } + }; + config.withConfiguration().listener(loggingChangeStateListener); + } +}