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 extends ExecutionCommon> 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"),