From 866c063b0d9c1fa1d40098ca12a2071c23d87a85 Mon Sep 17 00:00:00 2001 From: ialbert Date: Fri, 17 Nov 2023 13:29:26 +0300 Subject: [PATCH] http://jira.mfd.msk:8088/browse/CLS-590 --- .../session/stage/impl/FinishingSession.java | 81 +++++++++++++------ 1 file changed, 58 insertions(+), 23 deletions(-) diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/FinishingSession.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/FinishingSession.java index aabb6ae71..cd6176096 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/FinishingSession.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/FinishingSession.java @@ -28,20 +28,20 @@ import ru.spcex.platform.imdg.api.ImdgProvider; import ru.spcex.platform.imdg.api.predicate.ImdgPredicate; import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder; import ru.spcex.platform.imdg.api.predicate.specific.RegistryCodeSqlBuilder; +import ru.spcex.platform.utils.collection.Pair; import ru.spcex.platform.utils.enumeration.EnumMessage; import ru.spcex.platform.utils.enumeration.IEnumKey; import ru.spcex.platform.utils.enumeration.IMessageResolver; import java.time.Instant; import java.time.LocalDate; -import java.util.Collection; -import java.util.List; -import java.util.Objects; -import java.util.Set; +import java.util.*; +import java.util.function.Function; import java.util.stream.Collectors; import java.util.stream.Stream; import static ru.spcex.clearing.error.ClearingErrorInternal.SessionGeneralError; +import static ru.spcex.platform.utils.enumeration.IEnumKey.getEnumByKey; @Service @Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) @@ -109,48 +109,83 @@ public class FinishingSession implements ISessionStage { log.error("Session id={} not found. Can not update Execution's.", sessionId); } else { Collection oRegs = findORegistryBySessionId(sessionId); - Set exchangeExecutionIdAllowed = oRegs.stream() - .filter(rgs -> RegistryStatus.CLRD.equalsByKey(rgs.getRegistryStatus())) - .map(Registry::getGroupId) // Registry::getGroupId == Execution.ExchangeExecutionId, см. RegistryFondBuilder/RegistryDepoBuilder - .collect(Collectors.toSet()); + log.trace("found {} obligation registers by sessionId {}", oRegs.size(), sessionId); + 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-> { + RegistryStatus obligationStatus = obligationStatusByGroupId.get(id); + if (obligationStatus == null) return null; + return switch (obligationStatus) { + case CLRD -> CoverageStatus.ALWD; + case FAIL, UNCV, NACK, NACC -> CoverageStatus.DEND; + default -> null; + }; + }; Set exchangeExecutionIdsPreviousDay = oRegs.stream() .filter(rgs -> Objects.nonNull(rgs.getSettlementDate())) .filter(rgs -> Objects.nonNull(rgs.getTradingDate())) .filter(rgs -> rgs.getSettlementDate().compareTo(rgs.getTradingDate()) > 0) .map(Registry::getGroupId) .collect(Collectors.toSet()); - + log.trace("obligations number with second leg in the past: {}", exchangeExecutionIdsPreviousDay.size()); + int allowed = 0; + int notAllowed = 0; if (IEnumKey.contains(session.getSessionType(), SessionType.FINL, SessionType.MEDM)) { Collection executions = findExecutionDepositBySessionId(sessionId); - int notAllowed = 0; + log.trace("loaded {} ExecutionDeposits for sessionId {}", executions.size(), sessionId); for (ExecutionDeposit execution : executions) { - CoverageStatus toStatus = exchangeExecutionIdAllowed.contains(execution.getExchangeExecutionId()) ? - CoverageStatus.ALWD : CoverageStatus.DEND; - if (toStatus == CoverageStatus.DEND) + CoverageStatus toStatus = statusByExchangeId.apply(execution.getExchangeExecutionId()); + if (toStatus == null) { + log.trace("CoverageStatus not defined for execution.id={}", execution.getId()); + continue; + } + if (toStatus == CoverageStatus.DEND) { notAllowed++; + } + if (toStatus == CoverageStatus.ALWD) { + allowed++; + } execution.setCoverageStatus(toStatus.getKey()); execution.setUpdated(Instant.now()); executionDepositImdg.update(execution); } - log.trace("By sessionId={} update {} ExecutionDeposit: {} allowed, {} denied", sessionId, - executions.size(), exchangeExecutionIdAllowed.size(), notAllowed); + log.trace("By sessionId={} processed {} ExecutionDeposit: {} allowed, {} denied, skipped {}", + sessionId, + executions.size(), + allowed, + notAllowed, + executions.size() - (allowed + notAllowed)); } else if (IEnumKey.contains(session.getSessionType(), SessionType.TRDT, SessionType.IPOB, SessionType.IPO0, SessionType.IPOT)) { List executions = Stream.concat( findExecutionFondBySessionId(sessionId).stream(), - findExecutionFondBySessionId(exchangeExecutionIdsPreviousDay).stream()) + findExecutionFondByExchangeId(exchangeExecutionIdsPreviousDay).stream()) .toList(); - int notAllowed = 0; + log.trace("loaded {} ExecutionFonds for sessionId {}", executions.size(), sessionId); for (ExecutionFond execution : executions) { - CoverageStatus toStatus = exchangeExecutionIdAllowed.contains(execution.getExchangeExecutionId()) ? - CoverageStatus.ALWD : CoverageStatus.DEND; - if (toStatus == CoverageStatus.DEND) + CoverageStatus toStatus = statusByExchangeId.apply(execution.getExchangeExecutionId()); + if (toStatus == null) { + log.trace("CoverageStatus not defined for execution.id={}", execution.getId()); + continue; + } + if (toStatus == CoverageStatus.DEND) { notAllowed++; + } + if (toStatus == CoverageStatus.ALWD) { + allowed++; + } execution.setCoverageStatus(toStatus.getKey()); execution.setUpdated(Instant.now()); executionFondImdg.update(execution); } - log.trace("By sessionId={} update {} ExecutionFond: {} allowed, {} denied", sessionId, - executions.size(), exchangeExecutionIdAllowed.size(), notAllowed); + log.trace("By sessionId={} processed {} ExecutionFond: allowed {}, denied {}, skipped {}", + sessionId, + executions.size(), + allowed, + notAllowed, + executions.size() - (allowed + notAllowed)); } else { log.warn("Unsupported session[{}].SessionType={} for update Executions", sessionId, session.getSessionType()); } @@ -270,7 +305,7 @@ public class FinishingSession implements ISessionStage { return result; } - protected Collection findExecutionFondBySessionId(Set exchangeIds) { + protected Collection findExecutionFondByExchangeId(Set exchangeIds) { ImdgPredicateBuilder pb = executionFondImdg.predicateBuilder(); ImdgPredicate prdct = pb.and( pb.in("exchangeExecutionId", exchangeIds.toArray(new Long[0])),