ialbert 2023-11-17 13:29:26 +03:00
parent 911c90bc45
commit 866c063b0d

View file

@ -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<Registry> oRegs = findORegistryBySessionId(sessionId);
Set<Long> 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<Long, RegistryStatus> 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<Long, CoverageStatus> 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<Long> 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<ExecutionDeposit> 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<ExecutionFond> 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<ExecutionFond> findExecutionFondBySessionId(Set<Long> exchangeIds) {
protected Collection<ExecutionFond> findExecutionFondByExchangeId(Set<Long> exchangeIds) {
ImdgPredicateBuilder pb = executionFondImdg.predicateBuilder();
ImdgPredicate prdct = pb.and(
pb.in("exchangeExecutionId", exchangeIds.toArray(new Long[0])),