From d64d15dc17836f7e98904c639dd97b2c3f2a347b Mon Sep 17 00:00:00 2001 From: ialbert Date: Tue, 24 Jun 2025 18:14:39 +0300 Subject: [PATCH] http://jira.mfd.msk:8088/browse/CLS-872 exchangeExecutionId/groupId state machine --- .../action/EndStageNotificationAction.java | 23 +++++++------- .../state/action/FinishingSessionAction.java | 18 ++++++----- ...mingPaymentInstructionReturnMkrAction.java | 13 ++++---- .../state/action/InclusionToPoolAction.java | 17 +++++++---- ...pectionObligationsDepositReturnAction.java | 30 +++++++++++-------- .../action/InspectionObligationsV2Action.java | 27 +++++++++++------ .../action/ObligationAdmissionAction.java | 9 ++++-- ...equirementAndObligationCreationAction.java | 11 ++++--- 8 files changed, 90 insertions(+), 58 deletions(-) diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/EndStageNotificationAction.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/EndStageNotificationAction.java index f1348b1f2..8656bb58a 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/EndStageNotificationAction.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/EndStageNotificationAction.java @@ -2,7 +2,6 @@ package ru.spcex.clearing.session.state.action; import java.time.Instant; import java.util.Collection; -import java.util.Objects; import java.util.Optional; import java.util.stream.Collectors; import org.slf4j.Logger; @@ -14,6 +13,7 @@ import org.springframework.statemachine.StateContext; import org.springframework.stereotype.Service; import ru.clearing.classes.statics.data.misc.Session; import ru.clearing.classes.statics.data.registry.Registry; +import ru.spcex.clearing.component.GroupRgsKey; import static ru.spcex.clearing.error.ClearingErrorInternal.SessionGeneralError; import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.platform.messaging.domain.Consts; @@ -69,25 +69,24 @@ public class EndStageNotificationAction extends AbstractSessionActionForOkErrorH } Collection forRegistries = selectRegistry(); - Collection groups = forRegistries.stream() - .map(Registry::getGroupId) - .filter(Objects::nonNull) - .distinct().collect(Collectors.toList()); + Collection groups = forRegistries.stream() + .map(registry -> new GroupRgsKey(registry.getGroupId(), registry.getMarket())) + .distinct().collect(Collectors.toList()); log.debug("Sending notifications for {} groups ({} registers) on section {}", - groups, forRegistries.size(), section); - for (Long groupId : groups) { - log.trace("For registry group {} send notification", - groupId); + groups, forRegistries.size(), section); + for (GroupRgsKey groupKey : groups) { + log.trace("For registry group/market {}/{} send notification", + groupKey.groupId(), groupKey.market()); Optional sResult; if (Section.FOND.equals(section)) { - sResult = notificationDF14(groupId); + sResult = notificationDF14(groupKey.groupId()); } else if (Section.MKR.equals(section)) { - sResult = notificationDF05(groupId); + sResult = notificationDF05(groupKey.groupId()); } else { throw new IllegalArgumentException("Unsupported section " + section); } sResult.ifPresent(enumMessage -> - log.error("When sending groupId={} has error: {}", groupId, msgResolver.resolve(enumMessage))); + log.error("When sending groupId={} has error: {}", groupKey.groupId(), msgResolver.resolve(enumMessage))); } } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/FinishingSessionAction.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/FinishingSessionAction.java index 1224fcb75..7ee0b57f1 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/FinishingSessionAction.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/FinishingSessionAction.java @@ -26,6 +26,7 @@ import ru.clearing.classes.statics.data.execution.ExecutionFond; import ru.clearing.classes.statics.data.misc.Session; import ru.clearing.classes.statics.data.registry.Registry; import ru.clearing.classes.statics.data.sdf.SDf05; +import ru.spcex.clearing.component.GroupRgsKey; import static ru.spcex.clearing.error.ClearingErrorInternal.SessionGeneralError; import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.platform.messaging.domain.Consts; @@ -140,12 +141,13 @@ public class FinishingSessionAction extends AbstractSessionActionForOkErrorHandl log.error("Session id={} not found. Can not update Execution's.", sessionId); } else { Collection oRegs = findORegistryBySessionId(allLiabilitiesAndOKClaims); - Map obligationStatusByGroupId; + Map obligationStatusByGroupId; obligationStatusByGroupId = oRegs.stream() - .map(rgs -> new Pair<>(rgs.getGroupId(), getEnumByKey(RegistryStatus.class, rgs.getRegistryStatus()))) - .filter(pair -> pair.getSecond() != null) - .collect(Collectors.toMap(Pair::getFirst, Pair::getSecond, (o1, o2) -> o1)); - Function statusByExchangeId = id -> { + .map(rgs -> new Pair<>(new GroupRgsKey(rgs.getGroupId(), rgs.getMarket()), + getEnumByKey(RegistryStatus.class, rgs.getRegistryStatus()))) + .filter(pair -> pair.getSecond() != null) + .collect(Collectors.toMap(Pair::getFirst, Pair::getSecond, (o1, o2) -> o1)); + Function statusByExchangeId = id-> { RegistryStatus obligationStatus = obligationStatusByGroupId.get(id); if (obligationStatus == null) return null; return switch (obligationStatus) { @@ -169,7 +171,8 @@ public class FinishingSessionAction extends AbstractSessionActionForOkErrorHandl Collection executions = findExecutionDepositBySessionId(sessionId); log.trace("loaded {} ExecutionDeposits for sessionId {}", executions.size(), sessionId); for (ExecutionDeposit execution : executions) { - CoverageStatus toStatus = statusByExchangeId.apply(execution.getExchangeExecutionId()); + CoverageStatus toStatus = statusByExchangeId.apply( + new GroupRgsKey(execution.getExchangeExecutionId(), execution.getMarket())); if (toStatus == null) { log.trace("CoverageStatus not defined for execution.id={}", execution.getId()); continue; @@ -199,7 +202,8 @@ public class FinishingSessionAction extends AbstractSessionActionForOkErrorHandl Map execFond = new HashMap<>(); log.trace("loaded {} Executions for sessionId {}", executions.size(), sessionId); for (ExecutionCommon execution : executions) { - CoverageStatus toStatus = statusByExchangeId.apply(execution.getExchangeExecutionId()); + CoverageStatus toStatus = statusByExchangeId.apply( + new GroupRgsKey(execution.getExchangeExecutionId(), execution.getMarket())); if (toStatus == null) { log.trace("CoverageStatus not defined for execution.id={}", execution.getId()); continue; diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/FormingPaymentInstructionReturnMkrAction.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/FormingPaymentInstructionReturnMkrAction.java index 32adab9d5..013c2fd67 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/FormingPaymentInstructionReturnMkrAction.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/FormingPaymentInstructionReturnMkrAction.java @@ -18,6 +18,7 @@ import org.springframework.stereotype.Service; import ru.clearing.classes.statics.data.account.Account; import ru.clearing.classes.statics.data.payment.PaymentInstruction; import ru.clearing.classes.statics.data.registry.Registry; +import ru.spcex.clearing.component.GroupRgsKey; import ru.spcex.clearing.error.ClearingError; import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.service.builder.PaymentInstructionBuilderFinalMkrDeals; @@ -107,16 +108,18 @@ public class FormingPaymentInstructionReturnMkrAction extends AbstractSessionAct List pmtCreated = new ArrayList<>(); - Map> registriesByGroup = rgsAll - .stream(). - collect(Collectors.groupingBy(Registry::getGroupId)); + Map> registriesByGroup = rgsAll + .stream() + .collect(Collectors.groupingBy( + registry -> new GroupRgsKey(registry.getGroupId(), registry.getMarket()) + )); log.info("registry groups found {}", registriesByGroup.size()); - for (Map.Entry> entry : registriesByGroup.entrySet()) { + for (Map.Entry> entry : registriesByGroup.entrySet()) { List groupRgs = entry.getValue(); Optional lmtO = groupRgs.stream().filter(rgs -> equalByRgs(RegistryTradingParams.OM_T, rgs)).findFirst(); Optional cmtO = groupRgs.stream().filter(rgs -> equalByRgs(RegistryTradingParams.TM_T, rgs)).findFirst(); if (lmtO.isEmpty() || cmtO.isEmpty()) { - log.error("LM*T or CM*T not found for group {}", entry.getKey()); + log.error("LM*T or CM*T not found for group/market {}/{}", entry.getKey().groupId(), entry.getKey().market()); continue; } Registry lm_t = lmtO.get(); //obligation by money 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 f54b2d63e..e048fe24e 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 @@ -19,12 +19,14 @@ import org.springframework.beans.factory.config.ConfigurableBeanFactory; import org.springframework.context.annotation.Scope; 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.misc.Session; import ru.clearing.classes.statics.data.registry.Registry; +import ru.spcex.clearing.component.GroupRgsKey; import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.clearing.session.stage.TaskType; @@ -119,17 +121,17 @@ public class InclusionToPoolAction extends AbstractSessionActionForOkErrorHandli } log.info("Inclusion to pool with predicate: {}", registryPredicate.toString()); Collection obligations = registryImdg.getCollectionObjectsByPredicate(registryPredicate); - Map> registryByGroupId = obligations.stream() + Map> registryByGroupId = obligations.stream() .filter(registry -> registry.getSettlementDate().isEqual(LocalDate.now())) - .collect(Collectors.groupingBy(Registry::getGroupId)); + .collect(Collectors.groupingBy(registry -> new GroupRgsKey(registry.getGroupId(), registry.getMarket()))); Optional ssnType = Optional.of(sessionType); Map rgsToUpdate = new HashMap<>(); List execsToUpdate = new ArrayList<>(); boolean loadExecs = !IEnumKey.contains(sessionType, SessionType.TRDT, SessionType.CURR, SessionType.UNIT, SessionType.IPOT); - for (Map.Entry> entrySet : registryByGroupId.entrySet()) { - log.debug("Processing set of registry with groupId: {}", entrySet.getKey()); + 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()); @@ -152,7 +154,7 @@ public class InclusionToPoolAction extends AbstractSessionActionForOkErrorHandli } @SuppressWarnings("unchecked") - private Collection updateExecutions(Long sessionId, SessionType ssnTpe, Long rgsGroupId, String rgsSection) { + private Collection updateExecutions(Long sessionId, SessionType ssnTpe, GroupRgsKey groupRgsKey, String rgsSection) { Instant now = Instant.now(); if (ssnTpe == null || sessionId == null) { log.debug("will not update executions#sessionId - couldn't determine session type"); @@ -180,7 +182,10 @@ public class InclusionToPoolAction extends AbstractSessionActionForOkErrorHandli return Collections.emptyList(); } ImdgPredicateBuilder pb = execImdg.predicateBuilder(); - ImdgPredicate prdct = pb.equals("exchangeExecutionId", rgsGroupId); + ImdgPredicate prdct = pb.and( + pb.equals("exchangeExecutionId", groupRgsKey.groupId()), + StringUtils.hasText(groupRgsKey.market()) ? pb.equals("market", groupRgsKey.market()) : pb.alwaysTrue() + ); Collection execs = execImdg.getCollectionObjectsByPredicate(prdct); //Imdg finalExecImdg = execImdg; execs.forEach(e -> { diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/InspectionObligationsDepositReturnAction.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/InspectionObligationsDepositReturnAction.java index c4d7ebcc4..d32f84c1f 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/InspectionObligationsDepositReturnAction.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/InspectionObligationsDepositReturnAction.java @@ -18,6 +18,8 @@ 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.GroupRgsKeyComparator; import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.service.AssetTrio; import ru.spcex.clearing.service.integration.GatewayRequestCreator; @@ -56,6 +58,7 @@ public class InspectionObligationsDepositReturnAction extends AbstractSessionAct private final TradingTimeService tradingTimeService; private final GatewayRequester gateway; private final PlanBalanceCalc planBalanceCalc; + private static final GroupRgsKeyComparator groupRgsKeyComparator = new GroupRgsKeyComparator(); @Autowired public InspectionObligationsDepositReturnAction(ImdgProvider imdgProvider, RegistryManager rgsMng, AssetTBFProcessing assets, TradingTimeService tradingTimeService, GatewayRequester gateway, PlanBalanceCalc planBalanceCalc) { @@ -79,26 +82,27 @@ public class InspectionObligationsDepositReturnAction extends AbstractSessionAct Collection registriesToProcess = registryImdg.getCollectionObjectsBySQL(sqlCondition); Map rgssToStore = new HashMap<>(registriesToProcess.size()); //вторые ноги у сделок - нужна ли оптимизация? - List>> registriesByGroupSorted = registriesToProcess.stream() - .collect(Collectors.groupingBy(Registry::getGroupId)) - .entrySet() - .stream() - .sorted(Map.Entry.comparingByKey()) - .toList(); + List>> registriesByGroupSorted = registryImdg.getCollectionObjectsBySQL(sqlCondition) + .stream() + .collect(Collectors.groupingBy(registry -> new GroupRgsKey(registry.getGroupId(), registry.getMarket()))) + .entrySet() + .stream() + .sorted(groupRgsKeyComparator) + .toList(); log.info("registry groups found {}", registriesByGroupSorted.size()); - for (Map.Entry> entry : registriesByGroupSorted) { + for (Map.Entry> entry : registriesByGroupSorted) { List group = entry.getValue(); Optional omtInGroupO = group.stream().filter(registry -> equalByRgs(OM_T, registry)).findFirst(); Optional tmtInGroupO = group.stream().filter(registry -> equalByRgs(TM_T, registry)).findFirst(); if (omtInGroupO.isEmpty() || tmtInGroupO.isEmpty()) { - log.error("groupId {} failed to find OM*T/TM*T registry", entry.getKey()); + log.error("groupId/market {}/{} failed to find OM*T/TM*T registry", entry.getKey().groupId(), entry.getKey().market()); group.forEach(rgs -> updateRegistryStatus(rgs, registryStatusFailed(), rgssToStore)); // FAIL or MNG continue; } Registry omtRgs = omtInGroupO.get(); Registry tmtRgs = tmtInGroupO.get(); - log.debug("groupId={}, OM*T.id={}", entry.getKey(), omtRgs.getId()); + log.debug("groupId={}, market={}, OM*T.id={}", entry.getKey().groupId(), entry.getKey().market(), omtRgs.getId()); BigDecimal omtBalance = safeBD(omtRgs.getBalance()); //первая часть сделки депозита @@ -131,18 +135,20 @@ public class InspectionObligationsDepositReturnAction extends AbstractSessionAct Optional amfAssetO = searchAssetByOMT(omtRgs); Optional amfAssetReceiverO = searchAssetByOMT(tmtRgs); if (amfAssetO.isEmpty()) { - log.error("OM*T register.id={} groupId={} failed to find AM*F asset", omtRgs.getId(), entry.getKey()); + log.error("OM*T register.id={} groupId={}, market={} failed to find AM*F asset", omtRgs.getId(), + entry.getKey().groupId(), entry.getKey().market()); group.forEach(rgs -> updateRegistryStatus(rgs, registryStatusFailed(), rgssToStore)); // FAIL or MNG continue; } if (amfAssetReceiverO.isEmpty()) { - log.error("TM*T register.id={} groupId={} failed to find AM*F asset", tmtRgs.getId(), entry.getKey()); + log.error("TM*T register.id={} groupId={}, market={} failed to find AM*F asset", tmtRgs.getId(), + entry.getKey().groupId(), entry.getKey().market()); group.forEach(rgs -> updateRegistryStatus(rgs, registryStatusFailed(), rgssToStore)); // FAIL or MNG continue; } Registry amfAsset = amfAssetO.get(); - log.debug("groupId={}, AM*F.id={}", entry.getKey(), amfAsset.getId()); + log.debug("groupId={}, market={}, AM*F.id={}", entry.getKey().groupId(), entry.getKey().market(), amfAsset.getId()); BigDecimal amfBalance = safeBD(amfAsset.getBalance()); BiConsumer processAssetsAndSetSessionId = (amf, amount) -> { Optional asts = assets.processByAm_f(amf, amount, false); diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/InspectionObligationsV2Action.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/InspectionObligationsV2Action.java index 9abc82939..8b4d39875 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/InspectionObligationsV2Action.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/InspectionObligationsV2Action.java @@ -19,7 +19,9 @@ import org.springframework.beans.factory.config.ConfigurableBeanFactory; import org.springframework.context.annotation.Scope; import org.springframework.statemachine.StateContext; import org.springframework.stereotype.Service; +import org.springframework.util.StringUtils; import ru.clearing.classes.statics.data.registry.Registry; +import ru.spcex.clearing.component.GroupRgsKey; import ru.spcex.clearing.component.predicate.cash.registry.AssetByTcrCompanyAccountCashedPredicate; import ru.spcex.clearing.error.ClearingError; import ru.spcex.clearing.error.RgsError; @@ -122,8 +124,8 @@ public class InspectionObligationsV2Action extends AbstractSessionActionForOkErr Collection registriesToProcess = registryImdg.getCollectionObjectsBySQL(sqlCondition); Map rgssToStore = new HashMap<>(registriesToProcess.size()); //вторые ноги у сделок - нужна ли оптимизация? - List>> registriesByGroupSorted = registriesToProcess.stream() - .collect(Collectors.groupingBy(Registry::getGroupId)) + List>> registriesByGroupSorted = registriesToProcess.stream() + .collect(Collectors.groupingBy(registry -> new GroupRgsKey(registry.getGroupId(), registry.getMarket()))) .entrySet() .stream() .sorted((entry1, entry2) -> { @@ -144,15 +146,15 @@ public class InspectionObligationsV2Action extends AbstractSessionActionForOkErr log.info("found {} ({} groups) registries by sql: {}", registriesToProcess.size(), registriesByGroupSorted.size(), sqlCondition); ImdgPredicateBuilder rgsPb = registryImdg.predicateBuilder(); GROUP: - for (Map.Entry> entry : registriesByGroupSorted) { + for (Map.Entry> entry : registriesByGroupSorted) { List group = entry.getValue(); List obligationsInGroup = group.stream().filter(registry -> IEnumKey.getEnumByKey(RegistryDesignation.class, registry.getRegistryDesignation()) == RegistryDesignation.O).toList(); - log.debug("Find {} obligation with ids: {} in group: {}", obligationsInGroup.size(), + log.debug("Find {} obligation with ids: {} in group/market: {}/{}", obligationsInGroup.size(), obligationsInGroup.stream() .map(Registry::getId).collect(Collectors.toList()), - entry.getKey()); + entry.getKey().groupId(), entry.getKey().market()); boolean isUncovered = false; List checkResults = new ArrayList<>(); for (Registry obligation : obligationsInGroup) { @@ -185,7 +187,7 @@ public class InspectionObligationsV2Action extends AbstractSessionActionForOkErr }); RegistryInstrumentType mOrS = IEnumKey.getEnumByKey(RegistryInstrumentType.class, rgs.getRegistryInstrumentType()); if (mOrS == null) { - log.error("RegistryInstrumentType is null for groupId {} {}.id={}", entry.getKey(), + log.error("RegistryInstrumentType is null for groupId/market {}/{} {}.id={}", entry.getKey().groupId(), entry.getKey().market(), rgs.getRegistryCode(), rgs.getId()); failGroup.run(); continue; @@ -338,8 +340,9 @@ public class InspectionObligationsV2Action extends AbstractSessionActionForOkErr private Collection getRefundDateRgsIfPresent(List registries) { if (registries.stream().anyMatch(rgs -> rgs.getRefundDate() != null) && !SessionType.MEDM.equals(sessionType)) { - Optional> groupIdAndSettleDate = registries.stream() - .map(rgs -> new Pair<>(rgs.getGroupId(), rgs.getSettlementDate())) + //доп группировка по market не нужна, лист registries уже сгруппирован + Optional> groupIdAndSettleDate = registries.stream() + .map(rgs -> new Pair<>(new GroupRgsKey(rgs.getGroupId(), rgs.getMarket()), rgs.getSettlementDate())) .filter(pair -> pair.getFirst() != null) .filter(pair -> pair.getSecond() != null) .findFirst(); @@ -347,9 +350,15 @@ public class InspectionObligationsV2Action extends AbstractSessionActionForOkErr return Collections.emptyList(); } ImdgPredicateBuilder pb = registryImdg.predicateBuilder(); + String groupMarket = groupIdAndSettleDate.get().getFirst().market(); + ImdgPredicate marketCondition = StringUtils.hasText(groupMarket) ? + pb.or(pb.equals("market", groupMarket), pb.isNull("market")) : pb.alwaysTrue(); return registryImdg.getCollectionObjectsByPredicate( pb.and( - pb.equals("groupId", groupIdAndSettleDate.get().getFirst()), + pb.and( + pb.equals("groupId", groupIdAndSettleDate.get().getFirst().groupId()), + marketCondition + ), pb.not(pb.equals("settlementDate", groupIdAndSettleDate.get().getSecond())) ) ); diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/ObligationAdmissionAction.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/ObligationAdmissionAction.java index 920ef1057..7ee19d66a 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/ObligationAdmissionAction.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/ObligationAdmissionAction.java @@ -19,6 +19,7 @@ import ru.clearing.classes.statics.data.company.Company; import ru.clearing.classes.statics.data.company.relation.Relation; import ru.clearing.classes.statics.data.registry.Registry; import ru.clearing.classes.statics.data.registry.TradingClearingRegistry; +import ru.spcex.clearing.component.GroupRgsKey; import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.service.validation.admission.AccountActiveValidationRule; import ru.spcex.clearing.service.validation.admission.ClearingAvailableValidationRule; @@ -73,10 +74,11 @@ public class ObligationAdmissionAction extends AbstractSessionActionForOkErrorHa private void obligationAdmission(StateContext ctx) { Long sessionId = ctx.getExtendedState().get(DataEnum.sessionId, Long.class); Collection registries = registryImdg.getCollectionObjectsBySQL("sessionId = " + sessionId); - Map> byGroups = registries.stream().collect(Collectors.groupingBy(Registry::getGroupId)); + Map> byGroups = registries.stream() + .collect(Collectors.groupingBy(registry -> new GroupRgsKey(registry.getGroupId(), registry.getMarket()))); log.info("found {} ({} groups) registries with sessionId {}", registries.size(), byGroups.size(), sessionId); Map rgsToUpdate = new HashMap<>(); - for (Map.Entry> grpEntry : byGroups.entrySet()) { + for (Map.Entry> grpEntry : byGroups.entrySet()) { EnumMessage groupError = null; List rgsGroup = grpEntry.getValue(); for (Registry rgs : rgsGroup) { @@ -88,7 +90,8 @@ public class ObligationAdmissionAction extends AbstractSessionActionForOkErrorHa } } if (groupError != null) { - log.warn("error {} for registries groupId = {}", messageResolver.resolve(groupError), grpEntry.getKey()); + log.warn("error {} for registries groupId/market = {}/{}", messageResolver.resolve(groupError), + grpEntry.getKey().groupId(), grpEntry.getKey().market()); for (Registry rgs : rgsGroup) { rgs.setRegistryStatus(RegistryStatus.NACK.getKey()); rgsToUpdate.put(rgs.getId(), rgs); diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/RequirementAndObligationCreationAction.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/RequirementAndObligationCreationAction.java index c42149e8d..29f07d53e 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/RequirementAndObligationCreationAction.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/RequirementAndObligationCreationAction.java @@ -300,7 +300,8 @@ public class RequirementAndObligationCreationAction extends AbstractSessionActio exec.getCompanyId(), settlementDt(exec), exec.getSessionId(), - exec.getExchangeExecutionId() + exec.getExchangeExecutionId(), + exec.getMarket() ); } @@ -333,11 +334,12 @@ public class RequirementAndObligationCreationAction extends AbstractSessionActio exec.getCompanyId(), ((ExecutionDeposit) exec).getSecondLegSettlementDate(), exec.getSessionId(), - exec.getExchangeExecutionId()); + exec.getExchangeExecutionId(), + exec.getMarket()); } public Optional imdgRegSearch(RegistryTradingParams p, Long tcrId, Long companyId, LocalDate settlementDt, - Long sessionId, Long exchangeExecutionId) { + Long sessionId, Long exchangeExecutionId, String market) { String sql = RegistryCodeSqlBuilder.getInstance(p).build(); ImdgPredicateBuilder pb = registryImdg.predicateBuilder(); ImdgPredicate rgstrPredicate = pb.and(pb.sql(sql), @@ -345,7 +347,8 @@ public class RequirementAndObligationCreationAction extends AbstractSessionActio pb.equals("companyId", companyId), pb.equals("settlementDate", settlementDt), pb.equals("sessionId", sessionId), - pb.equals("groupId", exchangeExecutionId) + pb.equals("groupId", exchangeExecutionId), + pb.or(pb.equals("market", market), pb.isNull("market")) ); Registry rgs = registryImdg.getSingleObjectByPredicate(rgstrPredicate); log.trace("rgs.id={} found by: {}", rgs == null ? "'not found'" : rgs.getId(), rgstrPredicate.toString());