diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/component/predicate/cash/CategoryByCompanyIdCashingPredicate.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/component/predicate/cash/CategoryByCompanyIdCashingPredicate.java new file mode 100644 index 000000000..1852d0bb0 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/component/predicate/cash/CategoryByCompanyIdCashingPredicate.java @@ -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 cash) { + ImdgPredicate prdct = pb.equals("companyId", companyId); + return pb.cashed(prdct, cash, companyId); + } +} \ No newline at end of file diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/MarketCodesBySessionConfig.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/MarketCodesBySessionConfig.java index 2934b9f6f..4ff66d6af 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/MarketCodesBySessionConfig.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/MarketCodesBySessionConfig.java @@ -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> marketCodesForCK(ImdgProvider imdgProvider) { + return new MarketCodesProviderV2(imdgProvider, Section.CURR, MarketType.SCND, MarketType.SCSP); + } + @Bean(name = "marketCodesForT0Primary") public Supplier> 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> { + private final MarketType[] marketType2; + private volatile List 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 get() { + if (codes == null) { + synchronized (lock) { + if (codes == null) { + imdgProvider.waitAvailable(); + Imdg 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 markets = marketImdg.getCollectionObjectsByPredicate( + pb.and(prd) + ); + codes = markets.stream().map(Market::getCode).distinct().collect(Collectors.toList()); + } + } + } + return codes; + } + } } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_2/specific/PaymSessionStateMachineConfig.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_2/specific/PaymSessionStateMachineConfig.java index b6e034ace..9a92197d5 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_2/specific/PaymSessionStateMachineConfig.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_2/specific/PaymSessionStateMachineConfig.java @@ -86,7 +86,7 @@ public class PaymSessionStateMachineConfig extends EnumStateMachineConfigurerAda FormingPaymentInstructionAssetsAction formingPaymentInstructionAssets, AgainReviseStage3Action againReviseStage3Action, FinishingSessionAction finishingSession, - @Qualifier("marketCodesForCurr") + @Qualifier("marketCodesForCK") Supplier> marketCodes, SaveRegistriesAndContinueAction obligationAdmissionContinueAction, SaveRegistriesAfterInspectionAndContinueAction saveRegistriesAfterInspectionAndContinueAction, diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_2/specific/PrecSessionStateMachineConfig.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_2/specific/PrecSessionStateMachineConfig.java index 32e3cc0fd..9f79cf235 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_2/specific/PrecSessionStateMachineConfig.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_2/specific/PrecSessionStateMachineConfig.java @@ -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> 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 diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_2/specific/PrepSessionStateMachineConfig.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_2/specific/PrepSessionStateMachineConfig.java index 725b6c2ea..4bea97e1b 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_2/specific/PrepSessionStateMachineConfig.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_2/specific/PrepSessionStateMachineConfig.java @@ -66,7 +66,7 @@ public class PrepSessionStateMachineConfig extends EnumStateMachineConfigurerAda @Qualifier("requirementAndObligationCreationAction") RequirementAndObligationCreationAction reqAndOblAction, ObligationAdmissionAction obligationsAdmission, - @Qualifier("marketCodesForCurr") + @Qualifier("marketCodesForCK") Supplier> marketCodes, SaveRegistriesAndContinueAction obligationAdmissionContinueAction, @Qualifier("pauseWorkflowStatusAction") diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_2/specific/PrpcSessionStateMachineConfig.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_2/specific/PrpcSessionStateMachineConfig.java new file mode 100644 index 000000000..dd3c656f8 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_2/specific/PrpcSessionStateMachineConfig.java @@ -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 { + private final static Logger log = LoggerFactory.getLogger(PrpcSessionStateMachineConfig.class); + private final Imdg 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 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> 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 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 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 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 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 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 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 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 config) throws Exception { + StateMachineListenerAdapter loggingChangeStateListener + = new MachineMonitoringListener(SESSION_TYPE.getKey()); + config + .withConfiguration() + .machineId(SESSION_TYPE.getKey()) + .listener(loggingChangeStateListener); + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_2/specific/RsltSessionStateMachineConfig.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_2/specific/RsltSessionStateMachineConfig.java index d1d94f1f8..d258ea635 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_2/specific/RsltSessionStateMachineConfig.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_2/specific/RsltSessionStateMachineConfig.java @@ -86,7 +86,7 @@ public class RsltSessionStateMachineConfig extends EnumStateMachineConfigurerAda FormingPaymentInstructionAssetsAction formingPaymentInstructionAssets, AgainReviseStage3Action againReviseStage3Action, FinishingSessionAction finishingSession, - @Qualifier("marketCodesForCurr") + @Qualifier("marketCodesForCK") Supplier> marketCodes, SaveRegistriesAndContinueAction obligationAdmissionContinueAction, SaveRegistriesAfterInspectionAndContinueAction saveRegistriesAfterInspectionAndContinueAction, diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/InspectionObligationsPrecAction.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/InspectionObligationsPrecAction.java index 332dd32c9..655485872 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/InspectionObligationsPrecAction.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/InspectionObligationsPrecAction.java @@ -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 registryImdg; + private final Imdg cmcImdg; private final ImdgId idGen; private final IMessageResolver msgResolver = new SimpleMessageResolver(); + private Function 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 cmcFilter) { + this.cmcFilter = cmcFilter; + } + private final RegistryCashAssetCash a___Cash = new RegistryCashAssetCash("[all assets]"); + private final CashV2 cmcCash = new CategoryByCompanyIdCash(); @Override public void actualExecute(StateContext 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 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 registriesToProcess = registryImdg.getCollectionObjectsBySQL(sqlCondition); List errs = new ArrayList<>(); Map rgssToStore = new HashMap<>(registriesToProcess.size()); //вторые ноги у сделок - нужна ли оптимизация? @@ -142,7 +163,7 @@ public class InspectionObligationsPrecAction extends ReviseStage1Action { }); Map assetsByTcrAndSec = new HashMap<>(); - for (Map.Entry> entry : registriesByGroupSorted) { + GROUP: for (Map.Entry> entry : registriesByGroupSorted) { List 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 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 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 { + public CategoryByCompanyIdCash() { + super("categoryCash"); + } + + @Override + protected Long extractKey(ClearingMemberCategory obj) { + return obj.getCompanyId(); + } + } } diff --git a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/RegistryTradingParams.java b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/RegistryTradingParams.java index 477bf94e9..cb3ac703d 100644 --- a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/RegistryTradingParams.java +++ b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/RegistryTradingParams.java @@ -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); } } diff --git a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/SessionType.java b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/SessionType.java index 097309fd7..c3dd77765 100644 --- a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/SessionType.java +++ b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/SessionType.java @@ -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) {