From b8a4b114fa961d5c85d8cf6f4b6bd87d9ca18e4a Mon Sep 17 00:00:00 2001 From: ialbert Date: Mon, 6 Jul 2026 16:35:06 +0300 Subject: [PATCH] http://jira.mfd.msk:8088/browse/CLS-1060 IPOB + inclusion to pool action V2 --- ...aryAuctionBnSessionStateMachineConfig.java | 21 +- .../state/action/InclusionToPoolAction.java | 16 ++ .../state/action/InclusionToPoolActionV2.java | 251 ++++++++++++++++++ .../platform/enumeration/SessionType.java | 2 +- 4 files changed, 283 insertions(+), 7 deletions(-) create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/InclusionToPoolActionV2.java diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_2/specific/PrimaryAuctionBnSessionStateMachineConfig.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_2/specific/PrimaryAuctionBnSessionStateMachineConfig.java index cff72925b..0e255b568 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_2/specific/PrimaryAuctionBnSessionStateMachineConfig.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/state_machine_2/specific/PrimaryAuctionBnSessionStateMachineConfig.java @@ -27,7 +27,7 @@ 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.FormingPaymentInstructionSecuritiesAction; -import ru.spcex.clearing.session.state.action.InclusionToPoolAction; +import ru.spcex.clearing.session.state.action.InclusionToPoolActionV2; import ru.spcex.clearing.session.state.action.InspectionObligationsV2Action; import ru.spcex.clearing.session.state.action.ObligationAdmissionAction; import ru.spcex.clearing.session.state.action.RequirementAndObligationCreationAction; @@ -46,6 +46,7 @@ import ru.spcex.clearing.session.state.guard.StashedRegistriesPresentGuard; import ru.spcex.clearing.session.state.listener.MachineMonitoringListener; 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; @@ -63,7 +64,7 @@ public class PrimaryAuctionBnSessionStateMachineConfig private final DealsPrepareAction dealsPrepareAction; private final RequirementAndObligationCreationAction reqAndOblAction; private final ObligationAdmissionAction obligationAdmissionAction; - private final InclusionToPoolAction inclusionToPoolAction; + private final InclusionToPoolActionV2 inclusionToPoolAction; private final InspectionObligationsV2Action inspectionObligationsV2Action; private final FormingPaymentInstructionSecuritiesAction formingPaymentInstructionSecuritiesAction; private final FinishingSessionAction finishingSessionAction; @@ -96,7 +97,7 @@ public class PrimaryAuctionBnSessionStateMachineConfig @Qualifier("requirementAndObligationCreationAction") RequirementAndObligationCreationAction reqAndOblAction, ObligationAdmissionAction obligationAdmissionAction, - InclusionToPoolAction inclusionToPoolAction, + InclusionToPoolActionV2 inclusionToPoolAction, InspectionObligationsV2Action inspectionObligationsV2Action, FormingPaymentInstructionSecuritiesAction formingPaymentInstructionSecuritiesAction, FinishingSessionAction finishingSessionAction, @@ -129,10 +130,18 @@ public class PrimaryAuctionBnSessionStateMachineConfig //stages settings: this.dealsPrepareAction.searchForExecutions(ExecutionType.ExecutionFond); - ImdgPredicateBuilder execFondPb = executionFondImdg.predicateBuilder(); - dealsPrepareAction.addExecutionFondCondition(execFondPb.regex("settlementCode", "^B\\d{2}$")); - dealsPrepareAction.addExecutionFondCondition(execFondPb.in("market", marketCodes.get().toArray(new String[0]))); + ImdgPredicateBuilder pb = executionFondImdg.predicateBuilder(); + dealsPrepareAction.addExecutionFondCondition(pb.regex("settlementCode", "^B\\d{2}$")); + dealsPrepareAction.addExecutionFondCondition(pb.in("market", marketCodes.get().toArray(new String[0]))); this.obligationAdmissionContinueAction.setDataEnum(DataEnum.obligationAdmissionStashedRgs); + + this.inclusionToPoolAction.addRgsConditionSessionTypeIn( + SessionType.IPOB, SessionType.IPOM + ); + this.inclusionToPoolAction.addRgsConditionStatusIn( + RegistryStatus.PROC, RegistryStatus.MNG + ); + this.inclusionToPoolAction.setExecType(ExecutionType.ExecutionFond); } @Override diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/InclusionToPoolAction.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/InclusionToPoolAction.java index 2fc7a3d84..33320237d 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/InclusionToPoolAction.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/InclusionToPoolAction.java @@ -49,6 +49,22 @@ import ru.spcex.platform.utils.enumeration.IEnumKey; import ru.spcex.platform.utils.enumeration.IMessageResolver; import ru.spcex.platform.utils.text.TextUtil; +/** + * при изменениях в этом шаге необходимо перевести сессию на {@link InclusionToPoolActionV2}
+ * все настроечные методы для "повторения" логики из этого кода сделаны в новом классе.
+ * (можно добавить setExecTypeBySection для закрытия некоторых кейсов + * определения imdg execution'ов из updateExecutions) + * @see InclusionToPoolActionV2#addRegistryCondition + * @see InclusionToPoolActionV2#addRgsConditionSessionTypeIn + * @see InclusionToPoolActionV2#addRgsConditionSessionTypeEquals + * @see InclusionToPoolActionV2#addRgsConditionDynamic + * @see InclusionToPoolActionV2#addRgsConditionStatusIn + * @see InclusionToPoolActionV2#addRgsConditionStatusEquals + * @see InclusionToPoolActionV2#skipExecutionLoad + * @see InclusionToPoolActionV2#setExecType + * + */ +@Deprecated @Service @Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) public class InclusionToPoolAction extends AbstractSessionActionForOkErrorHandling { diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/InclusionToPoolActionV2.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/InclusionToPoolActionV2.java new file mode 100644 index 000000000..684a220da --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/InclusionToPoolActionV2.java @@ -0,0 +1,251 @@ +package ru.spcex.clearing.session.state.action; + +import java.time.Instant; +import java.time.LocalDate; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collection; +import java.util.Collections; +import java.util.Comparator; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.Objects; +import java.util.Optional; +import java.util.function.Function; +import java.util.stream.Collectors; +import java.util.stream.Stream; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.config.ConfigurableBeanFactory; +import org.springframework.context.annotation.Scope; +import org.springframework.statemachine.ExtendedState; +import org.springframework.statemachine.StateContext; +import org.springframework.stereotype.Service; +import org.springframework.util.StringUtils; +import ru.clearing.classes.statics.data.execution.ExecutionCommon; +import ru.clearing.classes.statics.data.execution.ExecutionCurrency; +import ru.clearing.classes.statics.data.execution.ExecutionDeposit; +import ru.clearing.classes.statics.data.execution.ExecutionFond; +import ru.clearing.classes.statics.data.registry.Registry; +import ru.spcex.clearing.component.GroupRgsKey; +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.util.ClearingUtil; +import ru.spcex.platform.classes.base.interfaces.ExecutionType; +import ru.spcex.platform.enumeration.RegistryDesignation; +import ru.spcex.platform.enumeration.RegistryInstrumentType; +import ru.spcex.platform.enumeration.RegistryStatus; +import ru.spcex.platform.enumeration.RegistryUnit; +import ru.spcex.platform.enumeration.SessionType; +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; +import ru.spcex.platform.utils.text.TextUtil; + +/** + * рефактор старого депрекейтед класса {@link InclusionToPoolAction} + */ +@Service +@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) +public class InclusionToPoolActionV2 extends AbstractSessionActionForOkErrorHandling { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final Imdg registryImdg; + private final Imdg executionDepositImdg; + private final Imdg executionFondImdg; + private final Imdg executionCurrImdg; + private SessionType sessionType; + private final ImdgPredicateBuilder rgsPb; + + private boolean loadExecs = true; + private ExecutionType execType; // + private final List registryConditionsAndV2 = new ArrayList<>(); + private final List> dynamicConditions = new ArrayList<>(); + + @Autowired + public InclusionToPoolActionV2(ImdgProvider imdgProvider) { + this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class); + this.executionDepositImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionDeposit, ExecutionDeposit.class); + this.executionFondImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionFond, ExecutionFond.class); + this.executionCurrImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionCurrency, ExecutionCurrency.class); + this.rgsPb = registryImdg.predicateBuilder(); + } + + public InclusionToPoolActionV2 addRegistryCondition(ImdgPredicate imdgPredicate) { + this.registryConditionsAndV2.add(imdgPredicate); + return this; + } + + public void addRgsConditionSessionTypeIn(SessionType... sessionTypes) { + ImdgPredicate sessionTypePredicate = rgsPb.or( + Arrays.stream(sessionTypes) + .map(sessionType -> rgsPb.equals("sessionType", sessionType.getKey())) + .toArray(ImdgPredicate[]::new) + ); + this.registryConditionsAndV2.add(sessionTypePredicate); + } + + /** + * можно использовать метод sessionInCondition, но не хочется лишнего or в предикате хазелкаста. + */ + public void addRgsConditionSessionTypeEquals(SessionType sessionType) { + ImdgPredicate sessionTypePredicate = rgsPb.equals( + "sessionType", sessionType.getKey() + ); + this.registryConditionsAndV2.add(sessionTypePredicate); + } + + public void addRgsConditionDynamic(Function dynamicCondition) { + this.dynamicConditions.add(dynamicCondition); + } + + public void addRgsConditionStatusIn(RegistryStatus... registryStatuses) { + ImdgPredicate registryStatusesCondition = rgsPb.or( + Arrays.stream(registryStatuses) + .map(registryStatus -> rgsPb.equals("registryStatus", registryStatus.getKey())) + .toArray(ImdgPredicate[]::new) + ); + this.registryConditionsAndV2.add(registryStatusesCondition); + } + + public void addRgsConditionStatusEquals(RegistryStatus registryStatus) { + ImdgPredicate registryStatusCondition = rgsPb.equals("registryStatus", registryStatus.getKey()); + this.registryConditionsAndV2.add(registryStatusCondition); + } + + public void skipExecutionLoad() { + this.loadExecs = false; + } + + public void setExecType(ExecutionType execType) { + this.execType = execType; + } + + @Override + protected void actualExecute(StateContext ctx) { + sessionType = ctx.getExtendedState().get(DataEnum.sessionType, SessionType.class); + Long sessionId = ctx.getExtendedState().get(DataEnum.sessionId, Long.class); + Long counterPartyId = ctx.getExtendedState().get(DataEnum.counterPartyId, Long.class); + SessionType sessionType = ctx.getExtendedState().get(DataEnum.sessionType, SessionType.class); + String sqlCondition = String.format("registryDesignation in ('%s', '%s') and " + + "registryInstrumentType in ('%s', '%s') and " + + "registryUnit = '%s'", + RegistryDesignation.O.getKey(), RegistryDesignation.T.getKey(), + RegistryInstrumentType.S.getKey(), RegistryInstrumentType.M.getKey(), + RegistryUnit.T.getKey()); + ImdgPredicateBuilder prdctBuilder = registryImdg.predicateBuilder(); + ImdgPredicate defaultCondition = prdctBuilder.sql(sqlCondition); + + ImdgPredicate registryPredicate; + ImdgPredicate[] conditions = Stream.of( + dynamicConditions.stream().map(d -> d.apply(ctx.getExtendedState())), + registryConditionsAndV2.stream(), + Stream.of(defaultCondition)) + .flatMap(Function.identity()) + .toArray(ImdgPredicate[]::new); + registryPredicate = prdctBuilder.and(conditions); + + log.info("Inclusion to pool with predicate: {}", registryPredicate.toString()); + Collection obligations = registryImdg.getCollectionObjectsByPredicate(registryPredicate); + Map> registryByGroupId = obligations.stream() + .filter(registry -> registry.getSettlementDate().isEqual(LocalDate.now())) + .collect(Collectors.groupingBy(registry -> new GroupRgsKey(registry.getGroupId(), registry.getMarket()))); + + Optional ssnType = Optional.of(sessionType); + Map rgsToUpdate = new HashMap<>(); + List execsToUpdate = new ArrayList<>(); + for (Map.Entry> entrySet : registryByGroupId.entrySet()) { + log.debug("Processing set of registry with groupId/market: {}/{}", entrySet.getKey().groupId(), entrySet.getKey().market()); + String rgsSection = null; + for (Registry registry : entrySet.getValue()) { + registry.setRegistryStatus(RegistryStatus.POOL.getKey()); + registry.setClearingDate(LocalDate.now()); + registry.setSessionId(sessionId); + ssnType.ifPresent(st -> registry.setSessionType(st.getKey())); + rgsToUpdate.put(registry.getId(), registry); + if (TextUtil.isEmpty(rgsSection)) { + rgsSection = registry.getSection(); + } + } + if (loadExecs) { + execsToUpdate.addAll(updateExecutionsV2(sessionId, entrySet.getKey(), rgsSection)); + } + } + registryImdg.putAll(rgsToUpdate); + if (loadExecs) { + updateExecsBatchV2(execsToUpdate); + } + } + + @SuppressWarnings("unchecked") + private Collection updateExecutionsV2(Long sessionId, GroupRgsKey groupRgsKey, String rgsSection) { + Instant now = Instant.now(); + if (sessionId == null) { + log.error("cannot update executions#sessionId - session id is null"); + return Collections.emptyList(); + } + Imdg execImdg = null; + if (execType == null) { + throw new IllegalStateException("execution type not set"); + } + switch (execType) { + case ExecutionDeposit -> execImdg = executionDepositImdg; + case ExecutionFond -> execImdg = executionFondImdg; + case ExecutionCurrency -> execImdg = executionCurrImdg; + } + if (execImdg == null) { + log.warn("couldn't define Execution Type for session {}. Will not update executions#sessionId", + this.sessionType); + return Collections.emptyList(); + } + ImdgPredicateBuilder pb = execImdg.predicateBuilder(); + ImdgPredicate prdct = pb.and( + pb.equals("exchangeExecutionId", groupRgsKey.groupId()), + StringUtils.hasText(groupRgsKey.market()) ? pb.equals("market", groupRgsKey.market()) : pb.alwaysTrue() + ); + Collection execs = (Collection) execImdg.getCollectionObjectsByPredicate(prdct); + execs.forEach(e -> { + e.setUpdated(now); + e.setSessionId(sessionId); + }); + return execs; + } + + private void updateExecsBatchV2(List execs) { + log.debug("updating {} executions", execs.size()); + execs.sort(Comparator.comparing(ExecutionCommon::type)); + log.debug("sorted executions by type"); + + Map m = new HashMap<>(); + for (int i = 0; i < execs.size(); i++) { + ExecutionCommon exec = execs.get(i); + ExecutionType execType = exec.type(); + log.debug("batch from {} position type {}", i, execType); + m.put(exec.getId(), exec); + int j = i + 1; + while (j < execs.size() && j < i + 100) { + ExecutionCommon other = execs.get(j); + if (!Objects.equals(other.type(), execType)) { + break; + } + m.put(other.getId(), other); + j++; + } + i = j; + log.debug("batch size {} type {}", m.size(), execType); + if (!m.isEmpty()) { + switch (execType) { + case ExecutionDeposit -> executionDepositImdg.putAll(ClearingUtil.castMap(m)); + case ExecutionFond -> executionFondImdg.putAll(ClearingUtil.castMap(m)); + case ExecutionCurrency -> executionCurrImdg.putAll(ClearingUtil.castMap(m)); + } + } + log.debug("batch size {} type {} done", m.size(), execType); + m.clear(); + } + } +} 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 c3dd77765..73c368465 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 @@ -3,7 +3,7 @@ package ru.spcex.platform.enumeration; import ru.spcex.platform.utils.enumeration.IEnumKey; public enum SessionType implements IEnumKey { - IPOB("IPOB"), IPO0("IPO0"), IPOT("IPOT"), TRDT("TRDT"), + IPOB("IPOB"), IPOM("IPOM"), IPO0("IPO0"), IPOT("IPOT"), TRDT("TRDT"), MEDM("MEDM"), FINL("FINL"), XDEP("XDEP"), LIQU("LIQU"), CURR("CURR"), UNIT("UNIT"), PREP("PREP"), PREC("PREC"), PAYM("PAYM"), RSLT("RSLT"),