clearing-service FinishingSession статусы для Execution. imdg добавил индексы
This commit is contained in:
parent
3ec26e9d91
commit
ef217e8ff7
4 changed files with 103 additions and 5 deletions
|
|
@ -6,6 +6,10 @@ import org.springframework.beans.factory.annotation.Autowired;
|
|||
import org.springframework.beans.factory.config.ConfigurableBeanFactory;
|
||||
import org.springframework.context.annotation.Scope;
|
||||
import org.springframework.stereotype.Service;
|
||||
import ru.clearing.classes.statics.data.execution.ExecutionCommon;
|
||||
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.clearing.classes.statics.data.sdf.SDf05;
|
||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||
|
|
@ -19,20 +23,22 @@ import ru.spcex.clearing.session.stage.ISessionStage;
|
|||
import ru.spcex.clearing.session.stage.StageResult;
|
||||
import ru.spcex.clearing.session.stage.Task;
|
||||
import ru.spcex.clearing.session.stage.task.FinishingSessionPayload;
|
||||
import ru.spcex.platform.enumeration.RegistryDesignation;
|
||||
import ru.spcex.platform.enumeration.RegistryStatus;
|
||||
import ru.spcex.platform.enumeration.RegistryTradingParams;
|
||||
import ru.spcex.platform.enumeration.Section;
|
||||
import ru.spcex.platform.enumeration.*;
|
||||
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.imdg.api.predicate.specific.RegistryCodeSqlBuilder;
|
||||
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.util.Collection;
|
||||
import java.util.Map;
|
||||
import java.util.Objects;
|
||||
import java.util.Set;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
import static ru.spcex.clearing.error.ClearingErrorInternal.SessionGeneralError;
|
||||
|
||||
|
|
@ -43,6 +49,9 @@ public class FinishingSession implements ISessionStage {
|
|||
|
||||
private final ImdgProvider imdgProvider;
|
||||
private final Imdg<Registry> registryImdg;
|
||||
private final Imdg<ExecutionFond> executionFondImdg;
|
||||
private final Imdg<ExecutionDeposit> executionDepositImdg;
|
||||
private final Imdg<Session> sessionImdg;
|
||||
private final Imdg<SDf05> sDf05Imdg;
|
||||
private final KafkaSender kafkaSender;
|
||||
private final IMessageResolver msgResolver;
|
||||
|
|
@ -53,7 +62,10 @@ public class FinishingSession implements ISessionStage {
|
|||
@Autowired
|
||||
public FinishingSession(ImdgProvider imdgProvider, KafkaSender kafkaSender, IMessageResolver msgResolver, Sdf05Sender sdf05Sender, Sdf14Sender sdf14Sender) {
|
||||
this.imdgProvider = imdgProvider;
|
||||
this.sessionImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Session, Session.class);
|
||||
this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
|
||||
this.executionFondImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionFond, ExecutionFond.class);
|
||||
this.executionDepositImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionDeposit, ExecutionDeposit.class);
|
||||
this.sDf05Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf05, SDf05.class);
|
||||
this.kafkaSender = kafkaSender;
|
||||
this.msgResolver = msgResolver;
|
||||
|
|
@ -86,7 +98,54 @@ public class FinishingSession implements ISessionStage {
|
|||
rgs.setRegistryStatus(RegistryStatus.CLRD.getKey());
|
||||
rgs.setUpdated(now);
|
||||
registryImdg.update(rgs);
|
||||
|
||||
});
|
||||
log.debug("For sessionId={} was updated {} registers for status={}", sessionId, claimsAndLiabilities.size(), RegistryStatus.CLRD.getKey());
|
||||
|
||||
// Отметка статуса в Execution*
|
||||
Session session = sessionImdg.getSingleObjectByID(sessionId);
|
||||
if (session == null) {
|
||||
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());
|
||||
|
||||
if (IEnumKey.contains(session.getSessionType(), SessionType.FINL, SessionType.MEDM)) {
|
||||
Collection<ExecutionDeposit> executions = findExecutionDepositBySessionId(sessionId);
|
||||
int notAllowed = 0;
|
||||
for (ExecutionDeposit execution : executions) {
|
||||
CoverageStatus toStatus = exchangeExecutionIdAllowed.contains(execution.getExchangeExecutionId()) ?
|
||||
CoverageStatus.ALWD : CoverageStatus.DNED;
|
||||
if (toStatus == CoverageStatus.DNED)
|
||||
notAllowed++;
|
||||
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);
|
||||
} else if (IEnumKey.contains(session.getSessionType(), SessionType.TRDT, SessionType.IPOB, SessionType.IPO0, SessionType.IPOT)) {
|
||||
Collection<ExecutionFond> executions = findExecutionFondBySessionId(sessionId);
|
||||
int notAllowed = 0;
|
||||
for (ExecutionFond execution : executions) {
|
||||
CoverageStatus toStatus = exchangeExecutionIdAllowed.contains(execution.getExchangeExecutionId()) ?
|
||||
CoverageStatus.ALWD : CoverageStatus.DNED;
|
||||
if (toStatus == CoverageStatus.DNED)
|
||||
notAllowed++;
|
||||
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);
|
||||
} else {
|
||||
log.warn("Unsupported session[{}].SessionType={} for update Executions", sessionId, session.getSessionType());
|
||||
}
|
||||
}
|
||||
|
||||
// 2. Отправка сообщений
|
||||
Collection<Registry> forRegistries = selectRegistry(sessionId);
|
||||
log.debug("Found {} registries for sessionId={}", forRegistries.size(), sessionId);
|
||||
|
|
@ -173,4 +232,33 @@ public class FinishingSession implements ISessionStage {
|
|||
log.debug("found {} claims and liabilities", result.size());
|
||||
return result;
|
||||
}
|
||||
|
||||
protected Collection<Registry> findORegistryBySessionId(Long sessionId) {
|
||||
ImdgPredicateBuilder pb = registryImdg.predicateBuilder();
|
||||
ImdgPredicate prdct = pb.and(
|
||||
pb.equals("sessionId", sessionId),
|
||||
// registryStatus - любой статус
|
||||
pb.equals("registryDesignation", RegistryDesignation.O.getKey()));
|
||||
Collection<Registry> result = registryImdg.getCollectionObjectsByPredicate(prdct);
|
||||
log.trace("found {} liabilities by query {} ", result.size(), prdct);
|
||||
result = result.stream()
|
||||
.filter(registry -> Objects.equals(registry.getValueDate(), registry.getSettlementDate()))
|
||||
.collect(Collectors.toList());
|
||||
log.debug("found {} liabilities 'O' by query {} and ValueDate==SettlementDate", result.size(), prdct);
|
||||
return result;
|
||||
}
|
||||
|
||||
protected Collection<ExecutionFond> findExecutionFondBySessionId(Long sessionId) {
|
||||
Collection<ExecutionFond> result = executionFondImdg.getCollectionObjectsByFieldValues(Map.of("sessionId", sessionId));
|
||||
log.trace("found {} ExecutionFond by sessionId={}", result.size(), sessionId);
|
||||
return result;
|
||||
}
|
||||
|
||||
protected Collection<ExecutionDeposit> findExecutionDepositBySessionId(Long sessionId) {
|
||||
Collection<ExecutionDeposit> result = executionDepositImdg.getCollectionObjectsByFieldValues(Map.of("sessionId", sessionId));
|
||||
log.trace("found {} ExecutionDeposit by sessionId={}", result.size(), sessionId);
|
||||
return result;
|
||||
}
|
||||
|
||||
|
||||
}
|
||||
|
|
|
|||
|
|
@ -35,6 +35,11 @@ public class ExecutionFondMapStore extends TemplateMapStore<ExecutionFond> {
|
|||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public String[] getIndexingField() {
|
||||
return new String[]{"exchangeExecutionId", "sessionId"};
|
||||
}
|
||||
|
||||
@Override
|
||||
public Iterable<Long> loadAllKeys() {
|
||||
return defaultLoadAllKeysOnTodayByField("SETTLEMENT_DATE", true);
|
||||
|
|
|
|||
|
|
@ -29,6 +29,11 @@ public class RegistryMapStore extends TemplateMapStore<Registry> {
|
|||
return "REGISTRY";
|
||||
}
|
||||
|
||||
@Override
|
||||
public String[] getIndexingField() {
|
||||
return new String[]{"sessionId"};
|
||||
}
|
||||
|
||||
@Override
|
||||
public String[] getFields() {
|
||||
return new String[]{
|
||||
|
|
|
|||
|
|
@ -28,7 +28,7 @@ public class ExecutionDepositMapStore extends TemplateMapStore<ExecutionDeposit>
|
|||
}
|
||||
|
||||
public String[] getIndexingField() {
|
||||
return new String[]{"exchangeExecutionId"};
|
||||
return new String[]{"exchangeExecutionId", "sessionId"};
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue