http://jira.mfd.msk:8088/browse/CLS-872 exchangeExecutionId/groupId state machine

This commit is contained in:
ialbert 2025-06-24 18:14:39 +03:00
parent b8d55c4627
commit d64d15dc17
8 changed files with 90 additions and 58 deletions

View file

@ -2,7 +2,6 @@ package ru.spcex.clearing.session.state.action;
import java.time.Instant; import java.time.Instant;
import java.util.Collection; import java.util.Collection;
import java.util.Objects;
import java.util.Optional; import java.util.Optional;
import java.util.stream.Collectors; import java.util.stream.Collectors;
import org.slf4j.Logger; import org.slf4j.Logger;
@ -14,6 +13,7 @@ import org.springframework.statemachine.StateContext;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.misc.Session; import ru.clearing.classes.statics.data.misc.Session;
import ru.clearing.classes.statics.data.registry.Registry; import ru.clearing.classes.statics.data.registry.Registry;
import ru.spcex.clearing.component.GroupRgsKey;
import static ru.spcex.clearing.error.ClearingErrorInternal.SessionGeneralError; import static ru.spcex.clearing.error.ClearingErrorInternal.SessionGeneralError;
import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.Consts; import ru.spcex.clearing.platform.messaging.domain.Consts;
@ -69,25 +69,24 @@ public class EndStageNotificationAction extends AbstractSessionActionForOkErrorH
} }
Collection<Registry> forRegistries = selectRegistry(); Collection<Registry> forRegistries = selectRegistry();
Collection<Long> groups = forRegistries.stream() Collection<GroupRgsKey> groups = forRegistries.stream()
.map(Registry::getGroupId) .map(registry -> new GroupRgsKey(registry.getGroupId(), registry.getMarket()))
.filter(Objects::nonNull) .distinct().collect(Collectors.toList());
.distinct().collect(Collectors.toList());
log.debug("Sending notifications for {} groups ({} registers) on section {}", log.debug("Sending notifications for {} groups ({} registers) on section {}",
groups, forRegistries.size(), section); groups, forRegistries.size(), section);
for (Long groupId : groups) { for (GroupRgsKey groupKey : groups) {
log.trace("For registry group {} send notification", log.trace("For registry group/market {}/{} send notification",
groupId); groupKey.groupId(), groupKey.market());
Optional<EnumMessage> sResult; Optional<EnumMessage> sResult;
if (Section.FOND.equals(section)) { if (Section.FOND.equals(section)) {
sResult = notificationDF14(groupId); sResult = notificationDF14(groupKey.groupId());
} else if (Section.MKR.equals(section)) { } else if (Section.MKR.equals(section)) {
sResult = notificationDF05(groupId); sResult = notificationDF05(groupKey.groupId());
} else { } else {
throw new IllegalArgumentException("Unsupported section " + section); throw new IllegalArgumentException("Unsupported section " + section);
} }
sResult.ifPresent(enumMessage -> 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)));
} }
} }

View file

@ -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.misc.Session;
import ru.clearing.classes.statics.data.registry.Registry; import ru.clearing.classes.statics.data.registry.Registry;
import ru.clearing.classes.statics.data.sdf.SDf05; import ru.clearing.classes.statics.data.sdf.SDf05;
import ru.spcex.clearing.component.GroupRgsKey;
import static ru.spcex.clearing.error.ClearingErrorInternal.SessionGeneralError; import static ru.spcex.clearing.error.ClearingErrorInternal.SessionGeneralError;
import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.Consts; 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); log.error("Session id={} not found. Can not update Execution's.", sessionId);
} else { } else {
Collection<Registry> oRegs = findORegistryBySessionId(allLiabilitiesAndOKClaims); Collection<Registry> oRegs = findORegistryBySessionId(allLiabilitiesAndOKClaims);
Map<Long, RegistryStatus> obligationStatusByGroupId; Map<GroupRgsKey, RegistryStatus> obligationStatusByGroupId;
obligationStatusByGroupId = oRegs.stream() obligationStatusByGroupId = oRegs.stream()
.map(rgs -> new Pair<>(rgs.getGroupId(), getEnumByKey(RegistryStatus.class, rgs.getRegistryStatus()))) .map(rgs -> new Pair<>(new GroupRgsKey(rgs.getGroupId(), rgs.getMarket()),
.filter(pair -> pair.getSecond() != null) getEnumByKey(RegistryStatus.class, rgs.getRegistryStatus())))
.collect(Collectors.toMap(Pair::getFirst, Pair::getSecond, (o1, o2) -> o1)); .filter(pair -> pair.getSecond() != null)
Function<Long, CoverageStatus> statusByExchangeId = id -> { .collect(Collectors.toMap(Pair::getFirst, Pair::getSecond, (o1, o2) -> o1));
Function<GroupRgsKey, CoverageStatus> statusByExchangeId = id-> {
RegistryStatus obligationStatus = obligationStatusByGroupId.get(id); RegistryStatus obligationStatus = obligationStatusByGroupId.get(id);
if (obligationStatus == null) return null; if (obligationStatus == null) return null;
return switch (obligationStatus) { return switch (obligationStatus) {
@ -169,7 +171,8 @@ public class FinishingSessionAction extends AbstractSessionActionForOkErrorHandl
Collection<ExecutionDeposit> executions = findExecutionDepositBySessionId(sessionId); Collection<ExecutionDeposit> executions = findExecutionDepositBySessionId(sessionId);
log.trace("loaded {} ExecutionDeposits for sessionId {}", executions.size(), sessionId); log.trace("loaded {} ExecutionDeposits for sessionId {}", executions.size(), sessionId);
for (ExecutionDeposit execution : executions) { for (ExecutionDeposit execution : executions) {
CoverageStatus toStatus = statusByExchangeId.apply(execution.getExchangeExecutionId()); CoverageStatus toStatus = statusByExchangeId.apply(
new GroupRgsKey(execution.getExchangeExecutionId(), execution.getMarket()));
if (toStatus == null) { if (toStatus == null) {
log.trace("CoverageStatus not defined for execution.id={}", execution.getId()); log.trace("CoverageStatus not defined for execution.id={}", execution.getId());
continue; continue;
@ -199,7 +202,8 @@ public class FinishingSessionAction extends AbstractSessionActionForOkErrorHandl
Map<Long, ExecutionFond> execFond = new HashMap<>(); Map<Long, ExecutionFond> execFond = new HashMap<>();
log.trace("loaded {} Executions for sessionId {}", executions.size(), sessionId); log.trace("loaded {} Executions for sessionId {}", executions.size(), sessionId);
for (ExecutionCommon execution : executions) { for (ExecutionCommon execution : executions) {
CoverageStatus toStatus = statusByExchangeId.apply(execution.getExchangeExecutionId()); CoverageStatus toStatus = statusByExchangeId.apply(
new GroupRgsKey(execution.getExchangeExecutionId(), execution.getMarket()));
if (toStatus == null) { if (toStatus == null) {
log.trace("CoverageStatus not defined for execution.id={}", execution.getId()); log.trace("CoverageStatus not defined for execution.id={}", execution.getId());
continue; continue;

View file

@ -18,6 +18,7 @@ import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.account.Account; import ru.clearing.classes.statics.data.account.Account;
import ru.clearing.classes.statics.data.payment.PaymentInstruction; import ru.clearing.classes.statics.data.payment.PaymentInstruction;
import ru.clearing.classes.statics.data.registry.Registry; import ru.clearing.classes.statics.data.registry.Registry;
import ru.spcex.clearing.component.GroupRgsKey;
import ru.spcex.clearing.error.ClearingError; import ru.spcex.clearing.error.ClearingError;
import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.service.builder.PaymentInstructionBuilderFinalMkrDeals; import ru.spcex.clearing.service.builder.PaymentInstructionBuilderFinalMkrDeals;
@ -107,16 +108,18 @@ public class FormingPaymentInstructionReturnMkrAction extends AbstractSessionAct
List<PaymentInstruction> pmtCreated = new ArrayList<>(); List<PaymentInstruction> pmtCreated = new ArrayList<>();
Map<Long, List<Registry>> registriesByGroup = rgsAll Map<GroupRgsKey, List<Registry>> registriesByGroup = rgsAll
.stream(). .stream()
collect(Collectors.groupingBy(Registry::getGroupId)); .collect(Collectors.groupingBy(
registry -> new GroupRgsKey(registry.getGroupId(), registry.getMarket())
));
log.info("registry groups found {}", registriesByGroup.size()); log.info("registry groups found {}", registriesByGroup.size());
for (Map.Entry<Long, List<Registry>> entry : registriesByGroup.entrySet()) { for (Map.Entry<GroupRgsKey, List<Registry>> entry : registriesByGroup.entrySet()) {
List<Registry> groupRgs = entry.getValue(); List<Registry> groupRgs = entry.getValue();
Optional<Registry> lmtO = groupRgs.stream().filter(rgs -> equalByRgs(RegistryTradingParams.OM_T, rgs)).findFirst(); Optional<Registry> lmtO = groupRgs.stream().filter(rgs -> equalByRgs(RegistryTradingParams.OM_T, rgs)).findFirst();
Optional<Registry> cmtO = groupRgs.stream().filter(rgs -> equalByRgs(RegistryTradingParams.TM_T, rgs)).findFirst(); Optional<Registry> cmtO = groupRgs.stream().filter(rgs -> equalByRgs(RegistryTradingParams.TM_T, rgs)).findFirst();
if (lmtO.isEmpty() || cmtO.isEmpty()) { 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; continue;
} }
Registry lm_t = lmtO.get(); //obligation by money Registry lm_t = lmtO.get(); //obligation by money

View file

@ -19,12 +19,14 @@ import org.springframework.beans.factory.config.ConfigurableBeanFactory;
import org.springframework.context.annotation.Scope; import org.springframework.context.annotation.Scope;
import org.springframework.statemachine.StateContext; import org.springframework.statemachine.StateContext;
import org.springframework.stereotype.Service; 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.ExecutionCommon;
import ru.clearing.classes.statics.data.execution.ExecutionCurrency; import ru.clearing.classes.statics.data.execution.ExecutionCurrency;
import ru.clearing.classes.statics.data.execution.ExecutionDeposit; import ru.clearing.classes.statics.data.execution.ExecutionDeposit;
import ru.clearing.classes.statics.data.execution.ExecutionFond; import ru.clearing.classes.statics.data.execution.ExecutionFond;
import ru.clearing.classes.statics.data.misc.Session; import ru.clearing.classes.statics.data.misc.Session;
import ru.clearing.classes.statics.data.registry.Registry; import ru.clearing.classes.statics.data.registry.Registry;
import ru.spcex.clearing.component.GroupRgsKey;
import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.clearing.session.stage.TaskType; 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()); log.info("Inclusion to pool with predicate: {}", registryPredicate.toString());
Collection<Registry> obligations = registryImdg.getCollectionObjectsByPredicate(registryPredicate); Collection<Registry> obligations = registryImdg.getCollectionObjectsByPredicate(registryPredicate);
Map<Long, List<Registry>> registryByGroupId = obligations.stream() Map<GroupRgsKey, List<Registry>> registryByGroupId = obligations.stream()
.filter(registry -> registry.getSettlementDate().isEqual(LocalDate.now())) .filter(registry -> registry.getSettlementDate().isEqual(LocalDate.now()))
.collect(Collectors.groupingBy(Registry::getGroupId)); .collect(Collectors.groupingBy(registry -> new GroupRgsKey(registry.getGroupId(), registry.getMarket())));
Optional<SessionType> ssnType = Optional.of(sessionType); Optional<SessionType> ssnType = Optional.of(sessionType);
Map<Long, Registry> rgsToUpdate = new HashMap<>(); Map<Long, Registry> rgsToUpdate = new HashMap<>();
List<ExecutionCommon> execsToUpdate = new ArrayList<>(); List<ExecutionCommon> execsToUpdate = new ArrayList<>();
boolean loadExecs = !IEnumKey.contains(sessionType, boolean loadExecs = !IEnumKey.contains(sessionType,
SessionType.TRDT, SessionType.CURR, SessionType.UNIT, SessionType.IPOT); SessionType.TRDT, SessionType.CURR, SessionType.UNIT, SessionType.IPOT);
for (Map.Entry<Long, List<Registry>> entrySet : registryByGroupId.entrySet()) { for (Map.Entry<GroupRgsKey, List<Registry>> entrySet : registryByGroupId.entrySet()) {
log.debug("Processing set of registry with groupId: {}", entrySet.getKey()); log.debug("Processing set of registry with groupId/market: {}/{}", entrySet.getKey().groupId(), entrySet.getKey().market());
String rgsSection = null; String rgsSection = null;
for (Registry registry : entrySet.getValue()) { for (Registry registry : entrySet.getValue()) {
registry.setRegistryStatus(RegistryStatus.POOL.getKey()); registry.setRegistryStatus(RegistryStatus.POOL.getKey());
@ -152,7 +154,7 @@ public class InclusionToPoolAction extends AbstractSessionActionForOkErrorHandli
} }
@SuppressWarnings("unchecked") @SuppressWarnings("unchecked")
private <T extends ExecutionCommon> Collection<ExecutionCommon> updateExecutions(Long sessionId, SessionType ssnTpe, Long rgsGroupId, String rgsSection) { private <T extends ExecutionCommon> Collection<ExecutionCommon> updateExecutions(Long sessionId, SessionType ssnTpe, GroupRgsKey groupRgsKey, String rgsSection) {
Instant now = Instant.now(); Instant now = Instant.now();
if (ssnTpe == null || sessionId == null) { if (ssnTpe == null || sessionId == null) {
log.debug("will not update executions#sessionId - couldn't determine session type"); log.debug("will not update executions#sessionId - couldn't determine session type");
@ -180,7 +182,10 @@ public class InclusionToPoolAction extends AbstractSessionActionForOkErrorHandli
return Collections.emptyList(); return Collections.emptyList();
} }
ImdgPredicateBuilder pb = execImdg.predicateBuilder(); 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<T> execs = execImdg.getCollectionObjectsByPredicate(prdct); Collection<T> execs = execImdg.getCollectionObjectsByPredicate(prdct);
//Imdg<T> finalExecImdg = execImdg; //Imdg<T> finalExecImdg = execImdg;
execs.forEach(e -> { execs.forEach(e -> {

View file

@ -18,6 +18,8 @@ import org.springframework.statemachine.StateContext;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.company.ClearingMemberCategory; import ru.clearing.classes.statics.data.company.ClearingMemberCategory;
import ru.clearing.classes.statics.data.registry.Registry; 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.imdg.IMDGDistributedNames;
import ru.spcex.clearing.service.AssetTrio; import ru.spcex.clearing.service.AssetTrio;
import ru.spcex.clearing.service.integration.GatewayRequestCreator; import ru.spcex.clearing.service.integration.GatewayRequestCreator;
@ -56,6 +58,7 @@ public class InspectionObligationsDepositReturnAction extends AbstractSessionAct
private final TradingTimeService tradingTimeService; private final TradingTimeService tradingTimeService;
private final GatewayRequester gateway; private final GatewayRequester gateway;
private final PlanBalanceCalc planBalanceCalc; private final PlanBalanceCalc planBalanceCalc;
private static final GroupRgsKeyComparator groupRgsKeyComparator = new GroupRgsKeyComparator();
@Autowired @Autowired
public InspectionObligationsDepositReturnAction(ImdgProvider imdgProvider, RegistryManager rgsMng, AssetTBFProcessing assets, TradingTimeService tradingTimeService, GatewayRequester gateway, PlanBalanceCalc planBalanceCalc) { public InspectionObligationsDepositReturnAction(ImdgProvider imdgProvider, RegistryManager rgsMng, AssetTBFProcessing assets, TradingTimeService tradingTimeService, GatewayRequester gateway, PlanBalanceCalc planBalanceCalc) {
@ -79,26 +82,27 @@ public class InspectionObligationsDepositReturnAction extends AbstractSessionAct
Collection<Registry> registriesToProcess = registryImdg.getCollectionObjectsBySQL(sqlCondition); Collection<Registry> registriesToProcess = registryImdg.getCollectionObjectsBySQL(sqlCondition);
Map<Long, Registry> rgssToStore = new HashMap<>(registriesToProcess.size()); //вторые ноги у сделок - нужна ли оптимизация? Map<Long, Registry> rgssToStore = new HashMap<>(registriesToProcess.size()); //вторые ноги у сделок - нужна ли оптимизация?
List<Map.Entry<Long, List<Registry>>> registriesByGroupSorted = registriesToProcess.stream() List<Map.Entry<GroupRgsKey, List<Registry>>> registriesByGroupSorted = registryImdg.getCollectionObjectsBySQL(sqlCondition)
.collect(Collectors.groupingBy(Registry::getGroupId)) .stream()
.entrySet() .collect(Collectors.groupingBy(registry -> new GroupRgsKey(registry.getGroupId(), registry.getMarket())))
.stream() .entrySet()
.sorted(Map.Entry.comparingByKey()) .stream()
.toList(); .sorted(groupRgsKeyComparator)
.toList();
log.info("registry groups found {}", registriesByGroupSorted.size()); log.info("registry groups found {}", registriesByGroupSorted.size());
for (Map.Entry<Long, List<Registry>> entry : registriesByGroupSorted) { for (Map.Entry<GroupRgsKey, List<Registry>> entry : registriesByGroupSorted) {
List<Registry> group = entry.getValue(); List<Registry> group = entry.getValue();
Optional<Registry> omtInGroupO = group.stream().filter(registry -> equalByRgs(OM_T, registry)).findFirst(); Optional<Registry> omtInGroupO = group.stream().filter(registry -> equalByRgs(OM_T, registry)).findFirst();
Optional<Registry> tmtInGroupO = group.stream().filter(registry -> equalByRgs(TM_T, registry)).findFirst(); Optional<Registry> tmtInGroupO = group.stream().filter(registry -> equalByRgs(TM_T, registry)).findFirst();
if (omtInGroupO.isEmpty() || tmtInGroupO.isEmpty()) { 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 group.forEach(rgs -> updateRegistryStatus(rgs, registryStatusFailed(), rgssToStore)); // FAIL or MNG
continue; continue;
} }
Registry omtRgs = omtInGroupO.get(); Registry omtRgs = omtInGroupO.get();
Registry tmtRgs = tmtInGroupO.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()); BigDecimal omtBalance = safeBD(omtRgs.getBalance());
//первая часть сделки депозита //первая часть сделки депозита
@ -131,18 +135,20 @@ public class InspectionObligationsDepositReturnAction extends AbstractSessionAct
Optional<Registry> amfAssetO = searchAssetByOMT(omtRgs); Optional<Registry> amfAssetO = searchAssetByOMT(omtRgs);
Optional<Registry> amfAssetReceiverO = searchAssetByOMT(tmtRgs); Optional<Registry> amfAssetReceiverO = searchAssetByOMT(tmtRgs);
if (amfAssetO.isEmpty()) { 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 group.forEach(rgs -> updateRegistryStatus(rgs, registryStatusFailed(), rgssToStore)); // FAIL or MNG
continue; continue;
} }
if (amfAssetReceiverO.isEmpty()) { 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 group.forEach(rgs -> updateRegistryStatus(rgs, registryStatusFailed(), rgssToStore)); // FAIL or MNG
continue; continue;
} }
Registry amfAsset = amfAssetO.get(); 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()); BigDecimal amfBalance = safeBD(amfAsset.getBalance());
BiConsumer<Registry, BigDecimal> processAssetsAndSetSessionId = (amf, amount) -> { BiConsumer<Registry, BigDecimal> processAssetsAndSetSessionId = (amf, amount) -> {
Optional<AssetTrio> asts = assets.processByAm_f(amf, amount, false); Optional<AssetTrio> asts = assets.processByAm_f(amf, amount, false);

View file

@ -19,7 +19,9 @@ import org.springframework.beans.factory.config.ConfigurableBeanFactory;
import org.springframework.context.annotation.Scope; import org.springframework.context.annotation.Scope;
import org.springframework.statemachine.StateContext; import org.springframework.statemachine.StateContext;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.springframework.util.StringUtils;
import ru.clearing.classes.statics.data.registry.Registry; 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.component.predicate.cash.registry.AssetByTcrCompanyAccountCashedPredicate;
import ru.spcex.clearing.error.ClearingError; import ru.spcex.clearing.error.ClearingError;
import ru.spcex.clearing.error.RgsError; import ru.spcex.clearing.error.RgsError;
@ -122,8 +124,8 @@ public class InspectionObligationsV2Action extends AbstractSessionActionForOkErr
Collection<Registry> registriesToProcess = registryImdg.getCollectionObjectsBySQL(sqlCondition); Collection<Registry> registriesToProcess = registryImdg.getCollectionObjectsBySQL(sqlCondition);
Map<Long, Registry> rgssToStore = new HashMap<>(registriesToProcess.size()); //вторые ноги у сделок - нужна ли оптимизация? Map<Long, Registry> rgssToStore = new HashMap<>(registriesToProcess.size()); //вторые ноги у сделок - нужна ли оптимизация?
List<Map.Entry<Long, List<Registry>>> registriesByGroupSorted = registriesToProcess.stream() List<Map.Entry<GroupRgsKey, List<Registry>>> registriesByGroupSorted = registriesToProcess.stream()
.collect(Collectors.groupingBy(Registry::getGroupId)) .collect(Collectors.groupingBy(registry -> new GroupRgsKey(registry.getGroupId(), registry.getMarket())))
.entrySet() .entrySet()
.stream() .stream()
.sorted((entry1, entry2) -> { .sorted((entry1, entry2) -> {
@ -144,15 +146,15 @@ public class InspectionObligationsV2Action extends AbstractSessionActionForOkErr
log.info("found {} ({} groups) registries by sql: {}", registriesToProcess.size(), registriesByGroupSorted.size(), sqlCondition); log.info("found {} ({} groups) registries by sql: {}", registriesToProcess.size(), registriesByGroupSorted.size(), sqlCondition);
ImdgPredicateBuilder rgsPb = registryImdg.predicateBuilder(); ImdgPredicateBuilder rgsPb = registryImdg.predicateBuilder();
GROUP: GROUP:
for (Map.Entry<Long, List<Registry>> entry : registriesByGroupSorted) { for (Map.Entry<GroupRgsKey, List<Registry>> entry : registriesByGroupSorted) {
List<Registry> group = entry.getValue(); List<Registry> group = entry.getValue();
List<Registry> obligationsInGroup = group.stream().filter(registry -> List<Registry> obligationsInGroup = group.stream().filter(registry ->
IEnumKey.getEnumByKey(RegistryDesignation.class, registry.getRegistryDesignation()) == RegistryDesignation.O).toList(); 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() obligationsInGroup.stream()
.map(Registry::getId).collect(Collectors.toList()), .map(Registry::getId).collect(Collectors.toList()),
entry.getKey()); entry.getKey().groupId(), entry.getKey().market());
boolean isUncovered = false; boolean isUncovered = false;
List<CheckResult> checkResults = new ArrayList<>(); List<CheckResult> checkResults = new ArrayList<>();
for (Registry obligation : obligationsInGroup) { for (Registry obligation : obligationsInGroup) {
@ -185,7 +187,7 @@ public class InspectionObligationsV2Action extends AbstractSessionActionForOkErr
}); });
RegistryInstrumentType mOrS = IEnumKey.getEnumByKey(RegistryInstrumentType.class, rgs.getRegistryInstrumentType()); RegistryInstrumentType mOrS = IEnumKey.getEnumByKey(RegistryInstrumentType.class, rgs.getRegistryInstrumentType());
if (mOrS == null) { 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()); rgs.getRegistryCode(), rgs.getId());
failGroup.run(); failGroup.run();
continue; continue;
@ -338,8 +340,9 @@ public class InspectionObligationsV2Action extends AbstractSessionActionForOkErr
private Collection<Registry> getRefundDateRgsIfPresent(List<Registry> registries) { private Collection<Registry> getRefundDateRgsIfPresent(List<Registry> registries) {
if (registries.stream().anyMatch(rgs -> rgs.getRefundDate() != null) && !SessionType.MEDM.equals(sessionType)) { if (registries.stream().anyMatch(rgs -> rgs.getRefundDate() != null) && !SessionType.MEDM.equals(sessionType)) {
Optional<Pair<Long, LocalDate>> groupIdAndSettleDate = registries.stream() //доп группировка по market не нужна, лист registries уже сгруппирован
.map(rgs -> new Pair<>(rgs.getGroupId(), rgs.getSettlementDate())) Optional<Pair<GroupRgsKey, LocalDate>> groupIdAndSettleDate = registries.stream()
.map(rgs -> new Pair<>(new GroupRgsKey(rgs.getGroupId(), rgs.getMarket()), rgs.getSettlementDate()))
.filter(pair -> pair.getFirst() != null) .filter(pair -> pair.getFirst() != null)
.filter(pair -> pair.getSecond() != null) .filter(pair -> pair.getSecond() != null)
.findFirst(); .findFirst();
@ -347,9 +350,15 @@ public class InspectionObligationsV2Action extends AbstractSessionActionForOkErr
return Collections.emptyList(); return Collections.emptyList();
} }
ImdgPredicateBuilder pb = registryImdg.predicateBuilder(); 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( return registryImdg.getCollectionObjectsByPredicate(
pb.and( 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())) pb.not(pb.equals("settlementDate", groupIdAndSettleDate.get().getSecond()))
) )
); );

View file

@ -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.company.relation.Relation;
import ru.clearing.classes.statics.data.registry.Registry; import ru.clearing.classes.statics.data.registry.Registry;
import ru.clearing.classes.statics.data.registry.TradingClearingRegistry; import ru.clearing.classes.statics.data.registry.TradingClearingRegistry;
import ru.spcex.clearing.component.GroupRgsKey;
import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.service.validation.admission.AccountActiveValidationRule; import ru.spcex.clearing.service.validation.admission.AccountActiveValidationRule;
import ru.spcex.clearing.service.validation.admission.ClearingAvailableValidationRule; import ru.spcex.clearing.service.validation.admission.ClearingAvailableValidationRule;
@ -73,10 +74,11 @@ public class ObligationAdmissionAction extends AbstractSessionActionForOkErrorHa
private void obligationAdmission(StateContext<TaskType, SsnEvent> ctx) { private void obligationAdmission(StateContext<TaskType, SsnEvent> ctx) {
Long sessionId = ctx.getExtendedState().get(DataEnum.sessionId, Long.class); Long sessionId = ctx.getExtendedState().get(DataEnum.sessionId, Long.class);
Collection<Registry> registries = registryImdg.getCollectionObjectsBySQL("sessionId = " + sessionId); Collection<Registry> registries = registryImdg.getCollectionObjectsBySQL("sessionId = " + sessionId);
Map<Long, List<Registry>> byGroups = registries.stream().collect(Collectors.groupingBy(Registry::getGroupId)); Map<GroupRgsKey, List<Registry>> 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); log.info("found {} ({} groups) registries with sessionId {}", registries.size(), byGroups.size(), sessionId);
Map<Long, Registry> rgsToUpdate = new HashMap<>(); Map<Long, Registry> rgsToUpdate = new HashMap<>();
for (Map.Entry<Long, List<Registry>> grpEntry : byGroups.entrySet()) { for (Map.Entry<GroupRgsKey, List<Registry>> grpEntry : byGroups.entrySet()) {
EnumMessage groupError = null; EnumMessage groupError = null;
List<Registry> rgsGroup = grpEntry.getValue(); List<Registry> rgsGroup = grpEntry.getValue();
for (Registry rgs : rgsGroup) { for (Registry rgs : rgsGroup) {
@ -88,7 +90,8 @@ public class ObligationAdmissionAction extends AbstractSessionActionForOkErrorHa
} }
} }
if (groupError != null) { 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) { for (Registry rgs : rgsGroup) {
rgs.setRegistryStatus(RegistryStatus.NACK.getKey()); rgs.setRegistryStatus(RegistryStatus.NACK.getKey());
rgsToUpdate.put(rgs.getId(), rgs); rgsToUpdate.put(rgs.getId(), rgs);

View file

@ -300,7 +300,8 @@ public class RequirementAndObligationCreationAction extends AbstractSessionActio
exec.getCompanyId(), exec.getCompanyId(),
settlementDt(exec), settlementDt(exec),
exec.getSessionId(), exec.getSessionId(),
exec.getExchangeExecutionId() exec.getExchangeExecutionId(),
exec.getMarket()
); );
} }
@ -333,11 +334,12 @@ public class RequirementAndObligationCreationAction extends AbstractSessionActio
exec.getCompanyId(), exec.getCompanyId(),
((ExecutionDeposit) exec).getSecondLegSettlementDate(), ((ExecutionDeposit) exec).getSecondLegSettlementDate(),
exec.getSessionId(), exec.getSessionId(),
exec.getExchangeExecutionId()); exec.getExchangeExecutionId(),
exec.getMarket());
} }
public Optional<Registry> imdgRegSearch(RegistryTradingParams p, Long tcrId, Long companyId, LocalDate settlementDt, public Optional<Registry> 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(); String sql = RegistryCodeSqlBuilder.getInstance(p).build();
ImdgPredicateBuilder pb = registryImdg.predicateBuilder(); ImdgPredicateBuilder pb = registryImdg.predicateBuilder();
ImdgPredicate rgstrPredicate = pb.and(pb.sql(sql), ImdgPredicate rgstrPredicate = pb.and(pb.sql(sql),
@ -345,7 +347,8 @@ public class RequirementAndObligationCreationAction extends AbstractSessionActio
pb.equals("companyId", companyId), pb.equals("companyId", companyId),
pb.equals("settlementDate", settlementDt), pb.equals("settlementDate", settlementDt),
pb.equals("sessionId", sessionId), 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); Registry rgs = registryImdg.getSingleObjectByPredicate(rgstrPredicate);
log.trace("rgs.id={} found by: {}", rgs == null ? "'not found'" : rgs.getId(), rgstrPredicate.toString()); log.trace("rgs.id={} found by: {}", rgs == null ? "'not found'" : rgs.getId(), rgstrPredicate.toString());