This commit is contained in:
parent
c75a4655c9
commit
930bf6152a
10 changed files with 428 additions and 15 deletions
|
|
@ -0,0 +1,19 @@
|
|||
package ru.spcex.clearing.component.predicate.cash;
|
||||
|
||||
import ru.clearing.classes.statics.data.company.ClearingMemberCategory;
|
||||
import ru.spcex.platform.imdg.api.predicate.ImdgPredicate;
|
||||
import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
|
||||
import ru.spcex.platform.imdg.iml.hazelcast.adapter.CashV2;
|
||||
|
||||
public class CategoryByCompanyIdCashingPredicate {
|
||||
private final Long companyId;
|
||||
|
||||
public CategoryByCompanyIdCashingPredicate(final Long companyId) {
|
||||
this.companyId = companyId;
|
||||
}
|
||||
|
||||
public ImdgPredicate cashed(ImdgPredicateBuilder pb, CashV2<Long, ClearingMemberCategory> cash) {
|
||||
ImdgPredicate prdct = pb.equals("companyId", companyId);
|
||||
return pb.cashed(prdct, cash, companyId);
|
||||
}
|
||||
}
|
||||
|
|
@ -13,6 +13,8 @@ import ru.spcex.platform.enumeration.MarketType;
|
|||
import ru.spcex.platform.enumeration.Section;
|
||||
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;
|
||||
|
||||
@Configuration
|
||||
public class MarketCodesBySessionConfig {
|
||||
|
|
@ -32,6 +34,11 @@ public class MarketCodesBySessionConfig {
|
|||
return new MarketCodesProvider(imdgProvider, Section.CURR, MarketType.SCND);
|
||||
}
|
||||
|
||||
@Bean(name = "marketCodesForCK")
|
||||
public Supplier<List<String>> marketCodesForCK(ImdgProvider imdgProvider) {
|
||||
return new MarketCodesProviderV2(imdgProvider, Section.CURR, MarketType.SCND, MarketType.SCSP);
|
||||
}
|
||||
|
||||
@Bean(name = "marketCodesForT0Primary")
|
||||
public Supplier<List<String>> marketCodesForT0Primary(ImdgProvider imdgProvider) {
|
||||
return new MarketCodesProvider(imdgProvider, Section.FOND, MarketType.PRMR);
|
||||
|
|
@ -67,4 +74,41 @@ public class MarketCodesBySessionConfig {
|
|||
return codes;
|
||||
}
|
||||
}
|
||||
|
||||
private static class MarketCodesProviderV2 implements Supplier<List<String>> {
|
||||
private final MarketType[] marketType2;
|
||||
private volatile List<String> codes;
|
||||
private final Object lock = new Object();
|
||||
private final ImdgProvider imdgProvider;
|
||||
private final Section section;
|
||||
|
||||
public MarketCodesProviderV2(ImdgProvider imdgProvider, Section section, MarketType... marketType) {
|
||||
this.imdgProvider = imdgProvider;
|
||||
this.section = section;
|
||||
this.marketType2 = marketType;
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<String> get() {
|
||||
if (codes == null) {
|
||||
synchronized (lock) {
|
||||
if (codes == null) {
|
||||
imdgProvider.waitAvailable();
|
||||
Imdg<Market> marketImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Market, Market.class);
|
||||
ImdgPredicateBuilder pb = marketImdg.predicateBuilder();
|
||||
ImdgPredicate[] prd = new ImdgPredicate[marketType2.length + 1];
|
||||
prd[0] = pb.equals("section", section.getKey());
|
||||
for (int i = 0; i < marketType2.length; i++) {
|
||||
prd[i + 1] = pb.equals("marketType", marketType2[i].getKey());
|
||||
}
|
||||
Collection<Market> markets = marketImdg.getCollectionObjectsByPredicate(
|
||||
pb.and(prd)
|
||||
);
|
||||
codes = markets.stream().map(Market::getCode).distinct().collect(Collectors.toList());
|
||||
}
|
||||
}
|
||||
}
|
||||
return codes;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -86,7 +86,7 @@ public class PaymSessionStateMachineConfig extends EnumStateMachineConfigurerAda
|
|||
FormingPaymentInstructionAssetsAction formingPaymentInstructionAssets,
|
||||
AgainReviseStage3Action againReviseStage3Action,
|
||||
FinishingSessionAction finishingSession,
|
||||
@Qualifier("marketCodesForCurr")
|
||||
@Qualifier("marketCodesForCK")
|
||||
Supplier<List<String>> marketCodes,
|
||||
SaveRegistriesAndContinueAction obligationAdmissionContinueAction,
|
||||
SaveRegistriesAfterInspectionAndContinueAction saveRegistriesAfterInspectionAndContinueAction,
|
||||
|
|
|
|||
|
|
@ -42,6 +42,7 @@ import ru.spcex.clearing.session.state.listener.MachineMonitoringListener;
|
|||
import static ru.spcex.clearing.util.StateMachineUtil.chain;
|
||||
import static ru.spcex.clearing.util.StateMachineUtil.invert;
|
||||
import ru.spcex.platform.classes.base.interfaces.ExecutionType;
|
||||
import ru.spcex.platform.enumeration.ClearingCategory;
|
||||
import ru.spcex.platform.enumeration.Section;
|
||||
import ru.spcex.platform.enumeration.SessionType;
|
||||
import ru.spcex.platform.imdg.api.Imdg;
|
||||
|
|
@ -90,7 +91,7 @@ public class PrecSessionStateMachineConfig extends EnumStateMachineConfigurerAda
|
|||
FormingPaymentInstructionAssetsAction formingPaymentInstructionAssets,
|
||||
AgainReviseStage3Action againReviseStage3Action,
|
||||
FinishingSessionAction finishingSession,
|
||||
@Qualifier("marketCodesForCurr")
|
||||
@Qualifier("marketCodesForCK")
|
||||
Supplier<List<String>> marketCodes,
|
||||
SaveRegistriesAndContinueAction obligationAdmissionContinueAction,
|
||||
SaveRegistriesAfterInspectionAndContinueAction saveRegistriesAfterInspectionAndContinueAction, SendNotificationAction sendNotificationAction,
|
||||
|
|
@ -112,10 +113,15 @@ public class PrecSessionStateMachineConfig extends EnumStateMachineConfigurerAda
|
|||
this.inspectionObligationsV2Action = inspectionObligations;
|
||||
this.obligationAdmissionContinueAction = obligationAdmissionContinueAction;
|
||||
|
||||
dealsPrepareAction.searchForExecutions(ExecutionType.ExecutionCurrency);
|
||||
this.dealsPrepareAction.searchForExecutions(ExecutionType.ExecutionCurrency);
|
||||
ImdgPredicateBuilder execFondPb = executionCurrencyImdg.predicateBuilder();
|
||||
dealsPrepareAction.addExecutionCurrencyCondition(execFondPb.in("market", marketCodes.get().toArray(new String[0])));
|
||||
this.dealsPrepareAction.addExecutionCurrencyCondition(
|
||||
execFondPb.in("market", marketCodes.get().toArray(new String[0]))
|
||||
);
|
||||
this.obligationAdmissionContinueAction.setDataEnum(DataEnum.obligationAdmissionStashedRgs);
|
||||
this.inspectionObligationsV2Action.setClearingMemberCategoryFilter(
|
||||
category -> !ClearingCategory.P.equalsByKey(category)
|
||||
);
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -66,7 +66,7 @@ public class PrepSessionStateMachineConfig extends EnumStateMachineConfigurerAda
|
|||
@Qualifier("requirementAndObligationCreationAction")
|
||||
RequirementAndObligationCreationAction reqAndOblAction,
|
||||
ObligationAdmissionAction obligationsAdmission,
|
||||
@Qualifier("marketCodesForCurr")
|
||||
@Qualifier("marketCodesForCK")
|
||||
Supplier<List<String>> marketCodes,
|
||||
SaveRegistriesAndContinueAction obligationAdmissionContinueAction,
|
||||
@Qualifier("pauseWorkflowStatusAction")
|
||||
|
|
|
|||
|
|
@ -0,0 +1,253 @@
|
|||
package ru.spcex.clearing.config.state_machine_2.specific;
|
||||
|
||||
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.Qualifier;
|
||||
import org.springframework.context.annotation.Configuration;
|
||||
import org.springframework.statemachine.action.Action;
|
||||
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 ru.clearing.classes.statics.data.execution.ExecutionCurrency;
|
||||
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.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.FinishingSessionAction;
|
||||
import ru.spcex.clearing.session.state.action.FormingPaymentInstructionAssetsAction;
|
||||
import ru.spcex.clearing.session.state.action.InclusionToPoolAction;
|
||||
import ru.spcex.clearing.session.state.action.InspectionObligationsPrecAction;
|
||||
import ru.spcex.clearing.session.state.action.ObligationAdmissionAction;
|
||||
import ru.spcex.clearing.session.state.action.RequirementAndObligationCreationAction;
|
||||
import ru.spcex.clearing.session.state.action.ReviseStage1Action;
|
||||
import ru.spcex.clearing.session.state.action.SaveRegistriesAfterInspectionAndContinueAction;
|
||||
import ru.spcex.clearing.session.state.action.SaveRegistriesAndContinueAction;
|
||||
import ru.spcex.clearing.session.state.action.SendNotificationAction;
|
||||
import ru.spcex.clearing.session.state.action.UpdateWorkflowStatusAction;
|
||||
import ru.spcex.clearing.session.state.guard.ErroneousRegistriesPresentGuard;
|
||||
import ru.spcex.clearing.session.state.guard.StashedRegistriesPresentGuard;
|
||||
import ru.spcex.clearing.session.state.listener.MachineMonitoringListener;
|
||||
import static ru.spcex.clearing.util.StateMachineUtil.chain;
|
||||
import static ru.spcex.clearing.util.StateMachineUtil.invert;
|
||||
import ru.spcex.platform.classes.base.interfaces.ExecutionType;
|
||||
import ru.spcex.platform.enumeration.ClearingCategory;
|
||||
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("prpcSessionStateMachineConfig")
|
||||
@EnableStateMachineFactory(name = "PRPC")
|
||||
public class PrpcSessionStateMachineConfig extends EnumStateMachineConfigurerAdapter<TaskType, SsnEvent> {
|
||||
private final static Logger log = LoggerFactory.getLogger(PrpcSessionStateMachineConfig.class);
|
||||
private final Imdg<Session> sessionImdg;
|
||||
private final ReviseStage1Action reviseStage1Action;
|
||||
private final DealsPrepareAction dealsPrepareAction;
|
||||
private final RequirementAndObligationCreationAction reqAndOblAction;
|
||||
private final ObligationAdmissionAction obligationAdmissionAction;
|
||||
private final InclusionToPoolAction inclusionToPoolAction;
|
||||
private final InspectionObligationsPrecAction inspectionObligationsV2Action;
|
||||
private final SaveRegistriesAndContinueAction obligationAdmissionContinueAction;
|
||||
private final SaveRegistriesAfterInspectionAndContinueAction saveRegistriesAfterInspectionAndContinueAction;
|
||||
private final SendNotificationAction sendNotificationAction;
|
||||
private final DiscardRegistriesAction discardOblAdmStash =
|
||||
new DiscardRegistriesAction(DataEnum.obligationAdmissionStashedRgs);
|
||||
private final DiscardRegistriesAction discardInspOblStash = new DiscardRegistriesAction(
|
||||
DataEnum.inspectionObligationStashedRgs
|
||||
);
|
||||
private final UpdateWorkflowStatusAction pauseWsAction;
|
||||
private final UpdateWorkflowStatusAction activeWsAction;
|
||||
|
||||
private final StashedRegistriesPresentGuard oblAdmGuard = new StashedRegistriesPresentGuard(
|
||||
DataEnum.obligationAdmissionStashedRgs
|
||||
);
|
||||
private final ErroneousRegistriesPresentGuard inspErrGuard = new ErroneousRegistriesPresentGuard();
|
||||
private final SessionType SESSION_TYPE = SessionType.PRPC;
|
||||
private final Section SECTION = Section.CURR;
|
||||
private static final Action<TaskType, SsnEvent> finishAction = context ->
|
||||
log.info("changing state to CL10; interceptor should update session status automatically.");
|
||||
|
||||
public PrpcSessionStateMachineConfig(ImdgProvider imdgProvider,
|
||||
ReviseStage1Action reviseStage1Action,
|
||||
DealsPrepareAction dealsPrepare,
|
||||
@Qualifier("requirementAndObligationCreationAction")
|
||||
RequirementAndObligationCreationAction reqAndOblAction,
|
||||
ObligationAdmissionAction obligationsAdmission,
|
||||
InclusionToPoolAction inclusionToPoolAction,
|
||||
InspectionObligationsPrecAction inspectionObligations,
|
||||
FormingPaymentInstructionAssetsAction formingPaymentInstructionAssets,
|
||||
AgainReviseStage3Action againReviseStage3Action,
|
||||
FinishingSessionAction finishingSession,
|
||||
@Qualifier("marketCodesForCK")
|
||||
Supplier<List<String>> marketCodes,
|
||||
SaveRegistriesAndContinueAction obligationAdmissionContinueAction,
|
||||
SaveRegistriesAfterInspectionAndContinueAction saveRegistriesAfterInspectionAndContinueAction, SendNotificationAction sendNotificationAction,
|
||||
@Qualifier("pauseWorkflowStatusAction")
|
||||
UpdateWorkflowStatusAction pauseWsAction,
|
||||
@Qualifier("activeWorkflowStatusAction")
|
||||
UpdateWorkflowStatusAction activeWsAction) {
|
||||
this.sessionImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Session, Session.class);
|
||||
this.saveRegistriesAfterInspectionAndContinueAction = saveRegistriesAfterInspectionAndContinueAction;
|
||||
this.sendNotificationAction = sendNotificationAction;
|
||||
this.pauseWsAction = pauseWsAction;
|
||||
this.activeWsAction = activeWsAction;
|
||||
Imdg<ExecutionCurrency> executionCurrencyImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionCurrency, ExecutionCurrency.class);
|
||||
this.reviseStage1Action = reviseStage1Action;
|
||||
this.dealsPrepareAction = dealsPrepare;
|
||||
this.reqAndOblAction = reqAndOblAction;
|
||||
this.obligationAdmissionAction = obligationsAdmission;
|
||||
this.inclusionToPoolAction = inclusionToPoolAction;
|
||||
this.inspectionObligationsV2Action = inspectionObligations;
|
||||
this.obligationAdmissionContinueAction = obligationAdmissionContinueAction;
|
||||
|
||||
this.dealsPrepareAction.searchForExecutions(ExecutionType.ExecutionCurrency);
|
||||
ImdgPredicateBuilder execFondPb = executionCurrencyImdg.predicateBuilder();
|
||||
this.dealsPrepareAction.addExecutionCurrencyCondition(
|
||||
execFondPb.in("market", marketCodes.get().toArray(new String[0]))
|
||||
);
|
||||
this.obligationAdmissionContinueAction.setDataEnum(DataEnum.obligationAdmissionStashedRgs);
|
||||
this.inspectionObligationsV2Action.setClearingMemberCategoryFilter(
|
||||
ClearingCategory.P::equalsByKey
|
||||
);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void configure(StateMachineStateConfigurer<TaskType, SsnEvent> states) throws Exception {
|
||||
states
|
||||
.withStates()
|
||||
.initial(TaskType.StartRevise, new CreateSessionAction(sessionImdg, SESSION_TYPE, SECTION))
|
||||
.states(new HashSet<>(Arrays.asList(
|
||||
TaskType.StartRevise,
|
||||
TaskType.StartRevisePart1,
|
||||
TaskType.DealsPrepare,
|
||||
TaskType.RequirementsAndObligationsCreate,
|
||||
TaskType.ObligationsAdmission,
|
||||
TaskType.InclusionToPool,
|
||||
TaskType.InspectionObligations,
|
||||
TaskType.FinishingSession,
|
||||
TaskType.PSEUDO_waitAfterObligationAdmissionError,
|
||||
TaskType.PSEUDO_waitAfterInspectionError,
|
||||
TaskType.PSEUDO_waitAfterObligationAdmissionRestore,
|
||||
TaskType.PSEUDO_waitAfterInspectionRestore
|
||||
)))
|
||||
.end(TaskType.FinishingSession);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void configure(StateMachineTransitionConfigurer<TaskType, SsnEvent> transitions) throws Exception {
|
||||
transitions
|
||||
.withExternal()
|
||||
.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.InclusionToPool)
|
||||
.action(inclusionToPoolAction)
|
||||
.and()
|
||||
.withExternal()
|
||||
.source(TaskType.InclusionToPool)
|
||||
.target(TaskType.ObligationsAdmission)
|
||||
.action(obligationAdmissionAction);
|
||||
|
||||
sourceObligationAdmission(transitions);
|
||||
sourcePseudoAfterOAPause(transitions);
|
||||
sourcePseudoAfterIOPause(transitions);
|
||||
sourceInspection(transitions);
|
||||
}
|
||||
|
||||
private void sourceObligationAdmission(StateMachineTransitionConfigurer<TaskType, SsnEvent> transitions) throws Exception {
|
||||
transitions
|
||||
.withExternal()
|
||||
.source(TaskType.ObligationsAdmission)
|
||||
.target(TaskType.PSEUDO_waitAfterObligationAdmissionError)
|
||||
.guard(oblAdmGuard)
|
||||
.action(pauseWsAction)
|
||||
.and()
|
||||
.withExternal()
|
||||
.source(TaskType.ObligationsAdmission).target(TaskType.InspectionObligations)
|
||||
.guard(invert(oblAdmGuard))
|
||||
.action(chain(obligationAdmissionContinueAction, inspectionObligationsV2Action));
|
||||
}
|
||||
|
||||
private void sourceInspection(StateMachineTransitionConfigurer<TaskType, SsnEvent> transitions) throws Exception {
|
||||
transitions
|
||||
.withExternal()
|
||||
.source(TaskType.InspectionObligations)
|
||||
.target(TaskType.PSEUDO_waitAfterInspectionError)
|
||||
.guard(inspErrGuard)
|
||||
.action(pauseWsAction)
|
||||
.and()
|
||||
.withExternal()
|
||||
.source(TaskType.InspectionObligations).target(TaskType.FinishingSession)
|
||||
.guard(invert(inspErrGuard))
|
||||
.action(saveRegistriesAfterInspectionAndContinueAction)
|
||||
.action(sendNotificationAction)
|
||||
.action(finishAction);
|
||||
}
|
||||
|
||||
private void sourcePseudoAfterOAPause(StateMachineTransitionConfigurer<TaskType, SsnEvent> transitions) throws Exception {
|
||||
transitions
|
||||
//continue and repeat after ObligationAdmission
|
||||
.withExternal()
|
||||
.source(TaskType.PSEUDO_waitAfterObligationAdmissionError).target(TaskType.InspectionObligations)
|
||||
.action(chain(activeWsAction, obligationAdmissionContinueAction, inspectionObligationsV2Action))
|
||||
.event(SsnEvent.CONTINUE);
|
||||
for (var taskType : List.of(TaskType.PSEUDO_waitAfterObligationAdmissionError, TaskType.PSEUDO_waitAfterObligationAdmissionRestore)) {
|
||||
transitions
|
||||
.withExternal()
|
||||
.source(taskType).target(TaskType.ObligationsAdmission)
|
||||
.action(chain(activeWsAction, discardOblAdmStash, obligationAdmissionAction))
|
||||
.event(SsnEvent.REPEAT);
|
||||
}
|
||||
}
|
||||
|
||||
private void sourcePseudoAfterIOPause(StateMachineTransitionConfigurer<TaskType, SsnEvent> transitions) throws Exception {
|
||||
transitions
|
||||
//continue and repeat after InspectionObligations
|
||||
.withExternal()
|
||||
.source(TaskType.PSEUDO_waitAfterInspectionError).target(TaskType.FinishingSession)
|
||||
.action(activeWsAction)
|
||||
.action(saveRegistriesAfterInspectionAndContinueAction)
|
||||
.action(sendNotificationAction)
|
||||
.action(finishAction)
|
||||
.event(SsnEvent.CONTINUE);
|
||||
for (var taskType : List.of(TaskType.PSEUDO_waitAfterInspectionError, TaskType.PSEUDO_waitAfterInspectionRestore)) {
|
||||
transitions
|
||||
.withExternal()
|
||||
.source(taskType).target(TaskType.InspectionObligations)
|
||||
.action(chain(activeWsAction, discardInspOblStash, inspectionObligationsV2Action))
|
||||
.event(SsnEvent.REPEAT);
|
||||
}
|
||||
}
|
||||
|
||||
@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);
|
||||
}
|
||||
}
|
||||
|
|
@ -86,7 +86,7 @@ public class RsltSessionStateMachineConfig extends EnumStateMachineConfigurerAda
|
|||
FormingPaymentInstructionAssetsAction formingPaymentInstructionAssets,
|
||||
AgainReviseStage3Action againReviseStage3Action,
|
||||
FinishingSessionAction finishingSession,
|
||||
@Qualifier("marketCodesForCurr")
|
||||
@Qualifier("marketCodesForCK")
|
||||
Supplier<List<String>> marketCodes,
|
||||
SaveRegistriesAndContinueAction obligationAdmissionContinueAction,
|
||||
SaveRegistriesAfterInspectionAndContinueAction saveRegistriesAfterInspectionAndContinueAction,
|
||||
|
|
|
|||
|
|
@ -2,6 +2,8 @@ package ru.spcex.clearing.session.state.action;
|
|||
|
||||
import java.math.BigDecimal;
|
||||
import java.time.Instant;
|
||||
import java.time.LocalDate;
|
||||
import java.time.LocalDateTime;
|
||||
import java.util.ArrayList;
|
||||
import java.util.Collection;
|
||||
import java.util.Comparator;
|
||||
|
|
@ -10,6 +12,7 @@ import java.util.List;
|
|||
import java.util.Map;
|
||||
import java.util.Optional;
|
||||
import java.util.function.Consumer;
|
||||
import java.util.function.Function;
|
||||
import java.util.stream.Collectors;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
|
@ -18,8 +21,10 @@ 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.component.GroupRgsKey;
|
||||
import ru.spcex.clearing.component.predicate.cash.CategoryByCompanyIdCashingPredicate;
|
||||
import ru.spcex.clearing.component.predicate.cash.registry.AssetByTcrCompanyAccountCashedPredicate;
|
||||
import ru.spcex.clearing.error.RgsError;
|
||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||
|
|
@ -42,19 +47,24 @@ 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.iml.hazelcast.adapter.CashV2;
|
||||
import static ru.spcex.platform.utils.enumeration.IEnumKey.getEnumByKey;
|
||||
import ru.spcex.platform.utils.enumeration.IMessageResolver;
|
||||
import ru.spcex.platform.utils.enumeration.SimpleMessageResolver;
|
||||
import ru.spcex.platform.utils.number.BigDecimalUtil;
|
||||
import static ru.spcex.platform.utils.number.BigDecimalUtil.safeBD;
|
||||
import ru.spcex.platform.utils.text.TextUtil;
|
||||
import ru.spcex.platform.utils.time.TimeUtil;
|
||||
|
||||
@Service
|
||||
@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE)
|
||||
public class InspectionObligationsPrecAction extends ReviseStage1Action {
|
||||
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||
private final Imdg<Registry> registryImdg;
|
||||
private final Imdg<ClearingMemberCategory> cmcImdg;
|
||||
private final ImdgId idGen;
|
||||
private final IMessageResolver msgResolver = new SimpleMessageResolver();
|
||||
private Function<String, Boolean> cmcFilter;
|
||||
private SessionType sessionType;
|
||||
private Long sessionId;
|
||||
|
||||
|
|
@ -62,20 +72,30 @@ public class InspectionObligationsPrecAction extends ReviseStage1Action {
|
|||
public InspectionObligationsPrecAction(ImdgProvider imdgProvider) {
|
||||
super(imdgProvider);
|
||||
this.registryImdg = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_Registry, Registry.class, null);
|
||||
this.cmcImdg = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_ClearingMemberCategory, ClearingMemberCategory.class, null);
|
||||
this.idGen = imdgProvider.getImdgIdGenerator();
|
||||
}
|
||||
|
||||
public void setClearingMemberCategoryFilter(Function<String, Boolean> cmcFilter) {
|
||||
this.cmcFilter = cmcFilter;
|
||||
}
|
||||
|
||||
private final RegistryCashAssetCash a___Cash = new RegistryCashAssetCash("[all assets]");
|
||||
private final CashV2<Long, ClearingMemberCategory> cmcCash = new CategoryByCompanyIdCash();
|
||||
|
||||
@Override
|
||||
public void actualExecute(StateContext<TaskType, SsnEvent> ctx) {
|
||||
try {
|
||||
super.actualExecute(ctx);
|
||||
if (cmcFilter == null) throw new IllegalStateException(
|
||||
"ClearingMemberCategory filter was not set"
|
||||
);
|
||||
this.sessionId = ctx.getExtendedState().get(DataEnum.sessionId, Long.class);
|
||||
this.sessionType = ctx.getExtendedState().get(DataEnum.sessionType, SessionType.class);
|
||||
inspectionObligations(ctx);
|
||||
} finally {
|
||||
a___Cash.clear();
|
||||
cmcCash.clear();
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -104,10 +124,11 @@ public class InspectionObligationsPrecAction extends ReviseStage1Action {
|
|||
|
||||
private void inspectionObligations(StateContext<TaskType, SsnEvent> ctx) {
|
||||
Instant now = Instant.now();
|
||||
zeroBalanceForPM_T(now);
|
||||
|
||||
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);
|
||||
List<RgsErr> errs = new ArrayList<>();
|
||||
Map<Long, Registry> rgssToStore = new HashMap<>(registriesToProcess.size()); //вторые ноги у сделок - нужна ли оптимизация?
|
||||
|
|
@ -142,7 +163,7 @@ public class InspectionObligationsPrecAction extends ReviseStage1Action {
|
|||
});
|
||||
|
||||
Map<PKey, RgsTrio> assetsByTcrAndSec = new HashMap<>();
|
||||
for (Map.Entry<GroupRgsKey, List<Registry>> entry : registriesByGroupSorted) {
|
||||
GROUP: for (Map.Entry<GroupRgsKey, List<Registry>> entry : registriesByGroupSorted) {
|
||||
List<Registry> group = entry.getValue();
|
||||
for (Registry rgs : group) {
|
||||
RegistryTradingParams rgsParams = paramsFromRegistry(rgs);
|
||||
|
|
@ -151,6 +172,21 @@ public class InspectionObligationsPrecAction extends ReviseStage1Action {
|
|||
rgs.getId(), rgs.getRegistryCode(), OM_T, TM_T);
|
||||
continue;
|
||||
}
|
||||
ImdgPredicate cmcPrdct = new CategoryByCompanyIdCashingPredicate(
|
||||
rgs.getCompanyId()
|
||||
).cashed(cmcImdg.predicateBuilder(), cmcCash);
|
||||
ClearingMemberCategory cmc = cmcImdg.getSingleObjectByPredicate(cmcPrdct);
|
||||
if (cmc == null || TextUtil.isEmpty(cmc.getClearingMemberCategory())) {
|
||||
log.error("couldn't find ClearingMemberCategory for group.id={}, rgs.id={}, company.id={}",
|
||||
entry.getKey().groupId(), rgs.getId(), rgs.getCompanyId());
|
||||
setGroupStatusError.accept(group);
|
||||
continue GROUP;
|
||||
}
|
||||
if (!cmcFilter.apply(cmc.getClearingMemberCategory())) {
|
||||
log.debug("group.id={}, rgs.id={}, company.id={}: category {}, skipping",
|
||||
entry.getKey().groupId(), rgs.getId(), rgs.getCompanyId(),cmc.getClearingMemberCategory());
|
||||
continue;
|
||||
}
|
||||
ImdgPredicate am_bPrdct = AssetByTcrCompanyAccountCashedPredicate
|
||||
.getPredicate(rgs, AM_B)
|
||||
.cashed(rgsPb, a___Cash);
|
||||
|
|
@ -165,7 +201,7 @@ public class InspectionObligationsPrecAction extends ReviseStage1Action {
|
|||
log.debug("groupId {}, {}.id={} - A**T not found. settings FAIL to group",
|
||||
entry.getKey(), rgs.getRegistryCode(), rgs.getId());
|
||||
setGroupStatusError.accept(group);
|
||||
continue;
|
||||
continue GROUP;
|
||||
}
|
||||
Registry am_b = Optional.ofNullable(
|
||||
registryImdg.getFirstObjectByPredicate(am_bPrdct)
|
||||
|
|
@ -205,15 +241,21 @@ public class InspectionObligationsPrecAction extends ReviseStage1Action {
|
|||
log.debug("tcr.id[{}], security.id[{}] - difference >= 0, skipping creating PM_T", key.tcrId, key.securityId);
|
||||
return;
|
||||
}
|
||||
Registry newPm_t = copyRgs(am_t, now, P___);
|
||||
newPm_t.setSessionId(sessionId);
|
||||
newPm_t.setSessionType(sessionType.getKey());
|
||||
newPm_t.setBalance(difference);
|
||||
rgssToStore.put(newPm_t.getId(), newPm_t);
|
||||
var pm_tO = searchPM_T(am_b.getTradingClearingRegistryId(),
|
||||
am_b.getCompanyId(),
|
||||
am_b.getAccountId(),
|
||||
am_b.getSecurityId(),
|
||||
TimeUtil.toDateTime(now));
|
||||
pm_tO.ifPresent(pmt -> pmt.setUpdated(now));
|
||||
Registry pmt = pm_tO.orElseGet(() -> copyRgs(am_t, now, P___));
|
||||
pmt.setSessionId(sessionId);
|
||||
pmt.setSessionType(sessionType.getKey());
|
||||
pmt.setBalance(difference);
|
||||
rgssToStore.put(pmt.getId(), pmt);
|
||||
if (notificationMessage[0] == null) {
|
||||
notificationMessage[0] = new StringBuilder("Сформированы PM*T регистры:");
|
||||
}
|
||||
notificationMessage[0].append("\nдля компании %s – регистр %d;".formatted(newPm_t.getShortName(), newPm_t.getId()));
|
||||
notificationMessage[0].append("\nдля компании %s – регистр %d;".formatted(pmt.getShortName(), pmt.getId()));
|
||||
});
|
||||
if (notificationMessage[0] != null) {
|
||||
ctx.getExtendedState().getVariables().put(DataEnum.notificationMessage, notificationMessage[0].toString());
|
||||
|
|
@ -222,6 +264,38 @@ public class InspectionObligationsPrecAction extends ReviseStage1Action {
|
|||
ctx.getExtendedState().getVariables().put(DataEnum.inspectionObligationStashedRgs, rgssToStore);
|
||||
}
|
||||
|
||||
private void zeroBalanceForPM_T(Instant now) {
|
||||
LocalDate nowD = TimeUtil.toLocalDate(now);
|
||||
ImdgPredicateBuilder pb = registryImdg.predicateBuilder();
|
||||
ImdgPredicate prdct = pb.and(
|
||||
pb.equals("registryDesignation", RegistryDesignation.P.getKey()),
|
||||
pb.equals("clearingDate", nowD),
|
||||
pb.equals("sessionType", sessionType.getKey())
|
||||
);
|
||||
Collection<Registry> rgss = registryImdg.getCollectionObjectsByPredicate(prdct);
|
||||
for (Registry r : rgss) {
|
||||
r.setBalance(BigDecimal.ZERO);
|
||||
r.setRegistryStatus(RegistryStatus.OK.getKey());
|
||||
r.setUpdated(now);
|
||||
registryImdg.update(r);
|
||||
}
|
||||
}
|
||||
|
||||
private Optional<Registry> searchPM_T(Long tcrId, Long cmpId, Long accId, Long secId, LocalDateTime today) {
|
||||
ImdgPredicateBuilder pb = registryImdg.predicateBuilder();
|
||||
ImdgPredicate prdct = pb.and(
|
||||
pb.sql(RegistryCodeSqlBuilder.getInstance(PM_T).build()),
|
||||
pb.equals("tradingClearingRegistryId", tcrId),
|
||||
pb.equals("companyId", cmpId),
|
||||
pb.equals("accountId", accId),
|
||||
pb.equals("securityId", secId),
|
||||
pb.equals("sessionType", sessionType.getKey()),
|
||||
pb.equals("clearingDate", today)
|
||||
);
|
||||
return Optional.ofNullable(registryImdg.getSingleObjectByPredicate(prdct));
|
||||
}
|
||||
|
||||
|
||||
private record PKey(Long tcrId, Long securityId) {
|
||||
}
|
||||
|
||||
|
|
@ -263,4 +337,15 @@ public class InspectionObligationsPrecAction extends ReviseStage1Action {
|
|||
log.debug("created {}.id={} by {}.id={}", newRgs.getRegistryCode(), newRgs.getId(), rgs.getRegistryCode(), rgs.getId());
|
||||
return newRgs;
|
||||
}
|
||||
|
||||
private static class CategoryByCompanyIdCash extends CashV2<Long, ClearingMemberCategory> {
|
||||
public CategoryByCompanyIdCash() {
|
||||
super("categoryCash");
|
||||
}
|
||||
|
||||
@Override
|
||||
protected Long extractKey(ClearingMemberCategory obj) {
|
||||
return obj.getCompanyId();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -132,6 +132,7 @@ public record RegistryTradingParams(RegistryDesignation registryDesignation,
|
|||
public final static RegistryTradingParams _S_T;
|
||||
public final static RegistryTradingParams AMAT;
|
||||
public final static RegistryTradingParams P___;
|
||||
public final static RegistryTradingParams PM_T;
|
||||
|
||||
static {
|
||||
OS_T = new RegistryTradingParams(RegistryDesignation.O,
|
||||
|
|
@ -287,6 +288,10 @@ public record RegistryTradingParams(RegistryDesignation registryDesignation,
|
|||
null,
|
||||
null,
|
||||
null);
|
||||
PM_T = new RegistryTradingParams(RegistryDesignation.P,
|
||||
RegistryInstrumentType.M,
|
||||
null,
|
||||
RegistryUnit.T);
|
||||
}
|
||||
|
||||
}
|
||||
|
|
|
|||
|
|
@ -7,6 +7,7 @@ public enum SessionType implements IEnumKey {
|
|||
MEDM("MEDM"), FINL("FINL"), XDEP("XDEP"), LIQU("LIQU"),
|
||||
CURR("CURR"), UNIT("UNIT"),
|
||||
PREP("PREP"), PREC("PREC"), PAYM("PAYM"), RSLT("RSLT"),
|
||||
PRPC("PRPC")
|
||||
;
|
||||
|
||||
SessionType(String key) {
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue