IntermediateMkrSession state config

This commit is contained in:
ialbert 2025-06-30 17:46:06 +03:00
parent b5f69e8b59
commit eb3577fd06
9 changed files with 352 additions and 6 deletions

View file

@ -5,6 +5,7 @@ import java.util.Arrays;
import java.util.EnumSet;
import java.util.HashSet;
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;
@ -86,6 +87,7 @@ public class FinalSessionStateMachineConfig
Sdf56Action sendSdf56Action,
ReviseStage1 reviseStage1Action,
DealsPrepareAction dealsPrepareAction,
@Qualifier("requirementAndObligationCreationAction")
RequirementAndObligationCreationAction reqAndOblAction,
ObligationAdmissionAction obligationAdmissionAction,
InclusionToPoolAction inclusionToPoolAction,

View file

@ -0,0 +1,266 @@
package ru.spcex.clearing.config.state_machine_2.specific;
import java.time.LocalDate;
import java.util.Arrays;
import java.util.EnumSet;
import java.util.HashSet;
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.guard.Guard;
import org.springframework.statemachine.listener.StateMachineListenerAdapter;
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.DataEnum;
import ru.spcex.clearing.session.state.SsnEvent;
import ru.spcex.clearing.session.state.action.AgainReviseStage3Action;
import ru.spcex.clearing.session.state.action.CreateSessionAction;
import ru.spcex.clearing.session.state.action.DealsPrepareAction;
import ru.spcex.clearing.session.state.action.DiscardRegistriesAction;
import ru.spcex.clearing.session.state.action.EndStageNotificationAction;
import ru.spcex.clearing.session.state.action.FinishingSessionAction;
import ru.spcex.clearing.session.state.action.FormingPaymentInstructionAssetsAction;
import ru.spcex.clearing.session.state.action.InclusionToPoolAction;
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.SaveRegistriesAndContinueAction;
import ru.spcex.clearing.session.state.action.Sdf56Action;
import ru.spcex.clearing.session.state.action.SdfReceivedAction;
import ru.spcex.clearing.session.state.guard.PaymentsWereCreatedGuard;
import ru.spcex.clearing.session.state.guard.PaymentsWereNotCreatedGuard;
import ru.spcex.clearing.session.state.guard.SdfGuard;
import ru.spcex.clearing.session.state.guard.SdfGuardExtractorAfterAgainRevise;
import ru.spcex.clearing.session.state.guard.SdfGuardExtractorAfterAssets;
import ru.spcex.clearing.session.state.guard.StashedRegistriesPresentGuard;
import ru.spcex.clearing.session.state.listener.MachineMonitoringListener;
import ru.spcex.clearing.util.StateMachineUtil;
import static ru.spcex.clearing.util.StateMachineUtil.*;
import ru.spcex.platform.classes.base.interfaces.ExecutionType;
import ru.spcex.platform.enumeration.RegistryStatus;
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.ImdgProvider;
import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
@Configuration("intermediateMkrSessionStateMachineFactoryConfig")
@EnableStateMachineFactory(name = "MEDM")
public class IntermediateMkrSessionStateMachineConfig
extends EnumStateMachineConfigurerAdapter<TaskType, SsnEvent> {
private final Imdg<Session> sessionImdg;
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 InspectionObligationsV2Action inspectionObligationsV2Action;
private final FormingPaymentInstructionAssetsAction formingPaymentInstructionAssetsAction;
private final AgainReviseStage3Action againRevise;
private final FinishingSessionAction finishingSessionAction;
private final EndStageNotificationAction endStageNotificationAction;
private final SaveRegistriesAndContinueAction obligationAdmissionContinueAction;
private final DiscardRegistriesAction discardOblAdmStash =
new DiscardRegistriesAction(DataEnum.obligationAdmissionStashedRgs);
private final StashedRegistriesPresentGuard oblAdmGuard = new StashedRegistriesPresentGuard(
DataEnum.obligationAdmissionStashedRgs
);
private final Guard<TaskType, SsnEvent> reviseSuccessGuard = guardCheckCtxForFlag(DataEnum.againReviseSuccess);
private final SessionType SESSION_TYPE = SessionType.MEDM;
private final Section SECTION = Section.MKR;
@Autowired
public IntermediateMkrSessionStateMachineConfig(
ImdgProvider imdgProvider,
Sdf56Action sendSdf56Action,
ReviseStage1 reviseStage1Action,
DealsPrepareAction dealsPrepareAction,
@Qualifier("requirementAndObligationCreationAction")
RequirementAndObligationCreationAction reqAndOblAction,
ObligationAdmissionAction obligationAdmissionAction,
InclusionToPoolAction inclusionToPoolAction,
InspectionObligationsV2Action inspectionObligationsV2Action,
FormingPaymentInstructionAssetsAction formingPaymentInstructionAssetsAction, AgainReviseStage3Action againRevise,
FinishingSessionAction finishingSessionAction, EndStageNotificationAction endStageNotificationAction, SaveRegistriesAndContinueAction obligationAdmissionContinueAction
) {
this.sendSdf56Action = sendSdf56Action;
this.sessionImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Session, Session.class);
this.againRevise = againRevise;
this.obligationAdmissionContinueAction = obligationAdmissionContinueAction;
Imdg<Registry> rgsImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
this.reviseStage1Action = reviseStage1Action;
this.dealsPrepareAction = dealsPrepareAction;
this.reqAndOblAction = reqAndOblAction;
this.obligationAdmissionAction = obligationAdmissionAction;
this.inclusionToPoolAction = inclusionToPoolAction;
this.inspectionObligationsV2Action = inspectionObligationsV2Action;
this.formingPaymentInstructionAssetsAction = formingPaymentInstructionAssetsAction;
this.finishingSessionAction = finishingSessionAction;
this.endStageNotificationAction = endStageNotificationAction;
this.finishingSessionAction.setSessionType(SESSION_TYPE);
this.finishingSessionAction.setSection(SECTION);
this.finishingSessionAction.setPr("4");
//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.equals("valueDate", LocalDate.now())
)
);
this.inspectionObligationsV2Action.setSessionType(SESSION_TYPE);
this.obligationAdmissionContinueAction.setDataEnum(DataEnum.obligationAdmissionStashedRgs);
}
@Override
public void configure(StateMachineStateConfigurer<TaskType, SsnEvent> states) throws Exception {
states
.withStates()
.initial(TaskType.StartRevise, chain(
new CreateSessionAction(sessionImdg)
.sessionType(SESSION_TYPE)
.section(SECTION),
sendSdf56Action))
.states(new HashSet<>(Arrays.asList(
TaskType.StartRevise,
TaskType.StartRevisePart1,
TaskType.DealsPrepare,
TaskType.RequirementsAndObligationsCreate,
TaskType.ObligationsAdmission,
TaskType.InclusionToPool,
TaskType.InspectionObligations,
TaskType.FormingRegistersOnOS,
TaskType.FormingPaymentInstruction,
TaskType.FinishingSession,
TaskType.EndStageNotification,
TaskType.PSEUDO_waitingForSdf,
TaskType.PSEUDO_waitingForObligationAdmission
)))
.end(TaskType.EndStageNotification);
}
@Override
public void configure(StateMachineTransitionConfigurer<TaskType, SsnEvent> transitions) throws Exception {
transitions
.withExternal()
.event(SsnEvent.SDF_57)
.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.PSEUDO_waitingForObligationAdmission)
.guard(oblAdmGuard)
.and()
.withExternal()
.source(TaskType.PSEUDO_waitingForObligationAdmission).target(TaskType.InclusionToPool)
.action(chain(obligationAdmissionContinueAction, inclusionToPoolAction))
.event(SsnEvent.CONTINUE)
.and()
.withExternal()
.source(TaskType.PSEUDO_waitingForObligationAdmission).target(TaskType.ObligationsAdmission)
.action(StateMachineUtil.chain(discardOblAdmStash, obligationAdmissionAction))
.event(SsnEvent.REPEAT)
.and()
.withExternal()
.source(TaskType.ObligationsAdmission).target(TaskType.InclusionToPool)
.guard(invert(oblAdmGuard))
.action(chain(obligationAdmissionContinueAction, inclusionToPoolAction))
.and()
.withExternal()
.source(TaskType.InclusionToPool).target(TaskType.InspectionObligations)
.action(inspectionObligationsV2Action)
.and()
.withExternal()
.source(TaskType.InspectionObligations).target(TaskType.FormingPaymentInstruction)
.action(formingPaymentInstructionAssetsAction)
.and()
//если не создалось paymentInstruction'ов
.withExternal()
.source(TaskType.FormingPaymentInstruction)
.target(TaskType.AgainRevise)
.guard(PaymentsWereNotCreatedGuard.instance)
.action(againRevise)
.and()
//если создались paymentInstruction, переходим в режим ожидания
.withExternal()
.source(TaskType.FormingPaymentInstruction)
.target(TaskType.PSEUDO_waitingForSdf)
.guard(PaymentsWereCreatedGuard.instance)
.and()
.withExternal()
.source(TaskType.PSEUDO_waitingForSdf)
.target(TaskType.AgainRevise)
.guard(new SdfGuard(SdfGuardExtractorAfterAssets.instance))
.action(againRevise)
.and()
.withExternal()
.source(TaskType.AgainRevise)
.target(TaskType.FinishingSession)
.action(finishingSessionAction)
.guard(reviseSuccessGuard)
.and()
.withExternal()
.source(TaskType.AgainRevise)
.target(TaskType.PSEUDO_waitingForSdfAfterAgainRevise)
.guard(invert(reviseSuccessGuard))
.and()
.withExternal()
.source(TaskType.PSEUDO_waitingForSdfAfterAgainRevise)
.target(TaskType.AgainRevise)
.guard(new SdfGuard(SdfGuardExtractorAfterAgainRevise.instance))
.action(againRevise)
.and()
.withExternal()
.source(TaskType.FinishingSession)
.target(TaskType.EndStageNotification)
.action(endStageNotificationAction);
for (var event: EnumSet.of(SsnEvent.SDF_01, SsnEvent.SDF_57, SsnEvent.SDF_04)) {
transitions
.withInternal()
.source(TaskType.PSEUDO_waitingForSdf)
.event(event)
.action(SdfReceivedAction.instance);
}
}
@Override
public void configure(StateMachineConfigurationConfigurer<TaskType, SsnEvent> config) throws Exception {
StateMachineListenerAdapter<TaskType, SsnEvent> loggingChangeStateListener
= new MachineMonitoringListener(SESSION_TYPE.getKey());
config
.withConfiguration()
.machineId(SESSION_TYPE.getKey())
.listener(loggingChangeStateListener);
}
}

View file

@ -57,6 +57,7 @@ public enum TaskType implements IEnumKey {
*/
EndStageNotification("CL11"),
PSEUDO_waitingForSdf("PSEUDO_waitingForSdf"),
PSEUDO_waitingForSdfAfterAgainRevise("PSEUDO_waitingForSdfAfterAgainRevise"),
PSEUDO_waitingForObligationAdmission("PSEUDO_waitingForObligationAdmission"),
;

View file

@ -15,4 +15,5 @@ public enum DataEnum {
paymentInstructionSecurity, //List<PaymentInstruction>
paymentsWereCreatedDepositReturn, //Collection<PaymentInstruction>
paymentInfo,
againReviseSuccess, //Boolean
}

View file

@ -9,13 +9,12 @@ 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.error.ClearingError;
import ru.spcex.clearing.error.ClearingRuntimeException;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.NotificationNewRequest;
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.platform.enumeration.AccountType;
import ru.spcex.platform.enumeration.ObjectType;
@ -26,7 +25,6 @@ import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.imdg.api.predicate.ImdgPredicate;
import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
import ru.spcex.platform.imdg.api.predicate.specific.RegistryCodeSqlBuilder;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import static ru.spcex.platform.utils.number.BigDecimalUtil.safeBD;
@Service
@ -75,13 +73,17 @@ public class AgainReviseStage3Action extends AbstractSessionActionForOkErrorHand
}
if (errorRegs > 0) {
ctx.getExtendedState().getVariables().put(DataEnum.againReviseSuccess, Boolean.FALSE);
log.warn("После сверки обнаружена разница между плановым и фактическим балансом. Всего {} регистров не совпали.", errorRegs);
log.warn("{} stage error", TaskType.AgainRevise);
NotificationNewRequest nRequest = new NotificationNewRequest();
nRequest.setObjectType(ObjectType.rgst.getKey());
nRequest.setPriority(Priority.HIGH.getKey());
nRequest.setComment("Сверка по результатам клиринговой сессии завершена с ошибками.");
kafkaSender.sendRequestToQueue(Consts.NOTIFICATION_NEW, nRequest);
throw new ClearingRuntimeException(new EnumMessage(ClearingError.PlanBalanceReviseError));
//throw new ClearingRuntimeException(new EnumMessage(ClearingError.PlanBalanceReviseError));
} else {
ctx.getExtendedState().getVariables().put(DataEnum.againReviseSuccess, Boolean.TRUE);
}
}
}

View file

@ -70,7 +70,7 @@ 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
@Service("requirementAndObligationCreationAction")
@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE)
public class RequirementAndObligationCreationAction extends AbstractSessionActionForOkErrorHandling {
private final Logger log = LoggerFactory.getLogger(getClass());

View file

@ -24,7 +24,7 @@ import ru.spcex.platform.classes.base.interfaces.ExecutionType;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.utils.collection.Pair;
@Service
@Service("requirementAndObligationCreationCompoundAction")
@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE)
public class RequirementAndObligationCreationCompoundAction extends RequirementAndObligationCreationAction {
private final Logger log = LoggerFactory.getLogger(getClass());

View file

@ -0,0 +1,66 @@
package ru.spcex.clearing.session.state.guard;
import java.util.EnumSet;
import java.util.Set;
import java.util.function.Function;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.statemachine.StateContext;
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.SdfTable;
import ru.spcex.platform.enumeration.Section;
public class SdfGuardExtractorAfterAgainRevise implements Function<StateContext<TaskType, SsnEvent>, Set<SdfTable>> {
private final Logger log = LoggerFactory.getLogger(getClass());
public static final SdfGuardExtractorAfterAgainRevise instance = new SdfGuardExtractorAfterAgainRevise();
@Override
public Set<SdfTable> apply(StateContext<TaskType, SsnEvent> ctx) {
Section section = ctx.getExtendedState().get(DataEnum.section, Section.class);
if (section == null) {
throw new IllegalStateException("section is null");
}
switch (section) {
case MKR -> {
return EnumSet.of(
SdfTable.SDF_01,
SdfTable.SDF_57
);
}
case FOND, MULT -> {
return EnumSet.of(
SdfTable.SDF_08,
SdfTable.SDF_21,
SdfTable.SDF_01,
SdfTable.SDF_57
);
}
default -> throw new IllegalStateException("unknown wait conditions for section " + section);
}
}
//public static SessionMonitor waitStep7(Section section) {
// switch (section) {
// case MKR -> {
// return SessionMonitor.create()
// .addCondition(new SdfCondition(SdfTable.SDF_04))
// .addCondition(new SdfCondition(SdfTable.SDF_01))
// .addCondition(new SdfCondition(SdfTable.SDF_57));
// }
// case FOND, MULT -> {
// return SessionMonitor.create()
// .addCondition(new SdfCondition(SdfTable.SDF_04))
// .addCondition(new SdfCondition(SdfTable.SDF_13))
// .addCondition(new SdfCondition(SdfTable.SDF_08))
// .addCondition(new SdfCondition(SdfTable.SDF_21))
// .addCondition(new SdfCondition(SdfTable.SDF_01))
// .addCondition(new SdfCondition(SdfTable.SDF_57));
// }
// default -> throw new IllegalStateException("unknown wait conditions for section " + section);
// }
//}
}

View file

@ -4,6 +4,7 @@ import org.springframework.statemachine.StateContext;
import org.springframework.statemachine.StateMachine;
import org.springframework.statemachine.action.Action;
import org.springframework.statemachine.guard.Guard;
import ru.spcex.clearing.session.state.DataEnum;
public class StateMachineUtil {
@ -30,4 +31,11 @@ public class StateMachineUtil {
};
}
public static <T1, T2> Guard<T1, T2> guardCheckCtxForFlag(DataEnum dataEnum) {
return context -> {
Boolean flag = (Boolean) context.getExtendedState().getVariables().get(dataEnum);
return flag != null && flag;
};
}
}