sessions step 3/4 optimization

imdg execution currency index
This commit is contained in:
ialbert 2024-07-31 13:56:34 +03:00
parent 21f5e6a680
commit 46da00357e
10 changed files with 342 additions and 151 deletions

View file

@ -0,0 +1,30 @@
package ru.spcex.clearing.component.predicate.cash.relation;
import ru.clearing.classes.statics.data.company.relation.Relation;
import ru.spcex.platform.imdg.api.predicate.ImdgPredicate;
import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
import ru.spcex.platform.imdg.iml.hazelcast.adapter.CashV2ByIdAndString;
public class RelationByConsumerIdAndServiceCashingPredicate {
private final CashV2ByIdAndString.CustomKey key;
private RelationByConsumerIdAndServiceCashingPredicate(Long consumerId, String service) {
this.key = new CashV2ByIdAndString.CustomKey(consumerId, service);
}
public static RelationByConsumerIdAndServiceCashingPredicate getPredicate(Long consumerId, String service) {
return new RelationByConsumerIdAndServiceCashingPredicate(consumerId, service);
}
public ImdgPredicate cashed(ImdgPredicateBuilder pb, CashV2ByIdAndString<Relation> cash) {
ImdgPredicate prdct = pb.and(
pb.equals("consumerId", key.id()),
pb.equals("service", key.code())
);
if (cash != null) {
return pb.cashed(prdct, cash, key);
} else {
return prdct;
}
}
}

View file

@ -54,7 +54,6 @@ import ru.spcex.clearing.service.validation.IdentificationFundsValidationRule;
import ru.spcex.clearing.service.validation.MarketIsUValidationRule;
import ru.spcex.clearing.service.validation.PaymentOutboundValidationRule;
import ru.spcex.clearing.service.validation.RefundDateValidationRule;
import ru.spcex.clearing.service.validation.RegistryStep3ValidationRule;
import ru.spcex.clearing.service.validation.ReturnDepositValidationRule;
import ru.spcex.clearing.service.validation.STradesValidationRule;
import ru.spcex.clearing.service.validation.Sdf01NewValidationRule;
@ -229,25 +228,6 @@ public class ValidationConfig {
};
}
@Bean("obligationAndRequirementsAdmissionValidator")
public Function<Registry, IValidator> registryValidator() {
return rgs -> {
ImdgValidationContext<Registry> context = new ImdgValidationContext<>();
context.setValidatedObject(rgs);
context.addImdg(IMDGDistributedNames.Map_Relation, imdgRelation);
context.addImdg(IMDGDistributedNames.Map_Session, imdgSession);
context.addImdg(IMDGDistributedNames.Map_SectionDictionary, imdgSectionDictionary);
context.addImdg(IMDGDistributedNames.Map_Account, imdgAccount);
context.addImdg(IMDGDistributedNames.Map_Company, imdgCompany);
context.addImdg(IMDGDistributedNames.Map_TradingClearingRegistry, imdgTradingClearingRegistry);
context.setLogPrefix(LogPrefixId.INSTANCE);
return new ValidatorImpl<>(context,
RegistryStep3ValidationRule.ClearingAvailable,
RegistryStep3ValidationRule.AccountActive,
RegistryStep3ValidationRule.CompanyActive,
RegistryStep3ValidationRule.TradingClearingRegistryActive);
};
}
@Bean("sdf01Validator")
public Function<SDf01, IValidator> sdf01Validator() {
return sDf01 -> {

View file

@ -1,107 +0,0 @@
package ru.spcex.clearing.service.validation;
import java.util.Map;
import java.util.Optional;
import ru.clearing.classes.statics.data.account.Account;
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.error.ClearingError;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.platform.enumeration.Section;
import ru.spcex.platform.enumeration.ServiceStatus;
import ru.spcex.platform.enumeration.WorkflowStatus;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.validation.ImdgValidationContext;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.enumeration.IEnumKey;
import ru.spcex.platform.utils.validation.IValidationRule;
public enum RegistryStep3ValidationRule implements IValidationRule<ImdgValidationContext<Registry>> {
ClearingAvailable() {
@Override
public Optional<EnumMessage> validate(ImdgValidationContext<Registry> context) {
Registry validatedObject = context.getValidatedObject();
// Supplier<String> sectionFind = () -> {
// Imdg<Session> sessionImdg = context.obtainMap(IMDGDistributedNames.Map_Session, Session.class);
// Imdg<SectionDictionary> sectionDictionaryImdg = context.obtainMap(IMDGDistributedNames.Map_SectionDictionary, SectionDictionary.class);
// Session session = sessionImdg.getSingleObjectByID(validatedObject.getSessionId());
// if (session == null) {
// return "null";
// } else {
// SectionDictionary section = sectionDictionaryImdg.getFirstObjectBySQL("code ='" + session.getSection() + "'");
// if (section == null) {
// return "null";
// } else {
// return section.getName();
// }
// }
// };
if (validatedObject.getCompanyId() == null) {
return of(ClearingError.ClearingUnavailableForCompany, validatedObject.getSection());
}
Imdg<Relation> relationImdg = context.obtainMap(IMDGDistributedNames.Map_Relation, Relation.class);
String section = validatedObject.getSection();
Section sctn = IEnumKey.getEnumByKey(Section.class, validatedObject.getSection());
if (Section.CURR.equals(sctn)) {
section = Section.MKR.getKey();
}
Relation relation = relationImdg.getFirstObjectByFieldValues(Map.of(
"consumerId", validatedObject.getCompanyId(),
"service", section));
if (relation == null || (!ServiceStatus.Active.equalsByKey(relation.getServiceStatus()) && !ServiceStatus.Reopened.equalsByKey(relation.getServiceStatus()))) {
return of(ClearingError.ClearingUnavailableForCompany, validatedObject.getSection(), validatedObject.getCompanyId());
}
context.storeObject(RegistryValidationStored.Relation, relation);
return empty();
}
},
AccountActive() {
@Override
public Optional<EnumMessage> validate(ImdgValidationContext<Registry> context) {
Registry validatedObject = context.getValidatedObject();
Imdg<Account> accountImdg = context.obtainMap(IMDGDistributedNames.Map_Account, Account.class);
Account account = accountImdg.getFirstObjectBySQL("id = %d".formatted(validatedObject.getAccountId()));
if (account == null || (!IEnumKey.contains(account.getStatus(), ServiceStatus.Active, ServiceStatus.Reopened))) {
return of(ClearingError.AccountNotActive, validatedObject.getAccountId());
}
return empty();
}
},
CompanyActive() {
@Override
public Optional<EnumMessage> validate(ImdgValidationContext<Registry> context) {
Registry validatedObject = context.getValidatedObject();
Imdg<Company> companyImdg = context.obtainMap(IMDGDistributedNames.Map_Company, Company.class);
if (validatedObject.getCompanyId() == null) {
return of(ClearingError.CompanyNotActive, validatedObject.getCompanyId());
}
Company company = companyImdg.getFirstObjectBySQL("id = %d".formatted(validatedObject.getCompanyId()));
if (company == null || !WorkflowStatus.Active.equalsByKey(company.getWorkflowStatus())) {
return of(ClearingError.CompanyNotActive, validatedObject.getCompanyId());
}
return empty();
}
},
TradingClearingRegistryActive() {
@Override
public Optional<EnumMessage> validate(ImdgValidationContext<Registry> context) {
Registry validatedObject = context.getValidatedObject();
if (validatedObject.getTradingClearingRegistryId() == null) {
return of(ClearingError.TradingClearingRegistryNotActive, validatedObject.getTradingClearingRegistryId());
}
Imdg<TradingClearingRegistry> tradingClearingRegistryImdg = context.obtainMap(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class);
TradingClearingRegistry tcr = tradingClearingRegistryImdg.getFirstObjectBySQL("id = " + validatedObject.getTradingClearingRegistryId());
if (tcr == null || !(ServiceStatus.Active.equalsByKey(tcr.getStatus()) || ServiceStatus.Reopened.equalsByKey(tcr.getStatus())))
return of(ClearingError.TradingClearingRegistryNotActive, validatedObject.getTradingClearingRegistryId());
return empty();
}
},
;
@Override
public String ruleName() {
return "RegistryStep3ValidationRule." + name();
}
}

View file

@ -0,0 +1,32 @@
package ru.spcex.clearing.service.validation.admission;
import java.util.Optional;
import ru.clearing.classes.statics.data.account.Account;
import ru.clearing.classes.statics.data.registry.Registry;
import ru.spcex.clearing.error.ClearingError;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.platform.enumeration.ServiceStatus;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.validation.ImdgValidationContext;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.enumeration.IEnumKey;
import ru.spcex.platform.utils.validation.IValidationRule;
public class AccountActiveValidationRule implements IValidationRule<ImdgValidationContext<Registry>> {
@Override
public String ruleName() {
return "AccountActive";
}
@Override
public Optional<EnumMessage> validate(ImdgValidationContext<Registry> context) {
Registry validatedObject = context.getValidatedObject();
Imdg<Account> accImdg = context.obtainMap(IMDGDistributedNames.Map_Account, Account.class);
Account account = accImdg.getSingleObjectByID(validatedObject.getAccountId());
if (account == null || (!IEnumKey.contains(account.getStatus(), ServiceStatus.Active, ServiceStatus.Reopened))) {
return of(ClearingError.AccountNotActive, validatedObject.getAccountId());
}
return empty();
}
}

View file

@ -0,0 +1,70 @@
package ru.spcex.clearing.service.validation.admission;
import java.util.Optional;
import ru.clearing.classes.statics.data.company.relation.Relation;
import ru.clearing.classes.statics.data.registry.Registry;
import ru.spcex.clearing.component.predicate.cash.relation.RelationByConsumerIdAndServiceCashingPredicate;
import ru.spcex.clearing.error.ClearingError;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.service.validation.RegistryValidationStored;
import ru.spcex.platform.enumeration.Section;
import ru.spcex.platform.enumeration.ServiceStatus;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.predicate.ImdgPredicate;
import ru.spcex.platform.imdg.iml.hazelcast.adapter.CashV2ByIdAndString;
import ru.spcex.platform.imdg.validation.ImdgValidationContext;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.enumeration.IEnumKey;
import ru.spcex.platform.utils.validation.IValidationRule;
public class ClearingAvailableValidationRule implements IValidationRule<ImdgValidationContext<Registry>> {
private final CashV2ByIdAndString<Relation> rltnCash;
public ClearingAvailableValidationRule(CashV2ByIdAndString<Relation> rltnCash) {
this.rltnCash = rltnCash;
}
@Override
public String ruleName() {
return "ClearingAvailable";
}
@Override
public Optional<EnumMessage> validate(ImdgValidationContext<Registry> context) {
Registry validatedObject = context.getValidatedObject();
//Supplier<String> sectionFind = () -> {
// Imdg<Session> sessionImdg = context.obtainMap(IMDGDistributedNames.Map_Session, Session.class);
// Imdg<SectionDictionary> sectionDictionaryImdg = context.obtainMap(IMDGDistributedNames.Map_SectionDictionary, SectionDictionary.class);
// Session session = sessionImdg.getSingleObjectByID(validatedObject.getSessionId());
// if (session == null) {
// return "null";
// } else {
// SectionDictionary section = sectionDictionaryImdg.getFirstObjectBySQL("code ='" + session.getSection() + "'");
// if (section == null) {
// return "null";
// } else {
// return section.getName();
// }
// }
//};
if (validatedObject.getCompanyId() == null) {
return of(ClearingError.ClearingUnavailableForCompany, validatedObject.getSection());
}
String section = validatedObject.getSection();
Section sctn = IEnumKey.getEnumByKey(Section.class, validatedObject.getSection());
if (Section.CURR.equals(sctn)) {
section = Section.MKR.getKey();
}
Imdg<Relation> rltnImdg = context.obtainMap(IMDGDistributedNames.Map_Relation, Relation.class);
ImdgPredicate prdct = RelationByConsumerIdAndServiceCashingPredicate
.getPredicate(validatedObject.getCompanyId(), section)
.cashed(rltnImdg.predicateBuilder(), rltnCash);
Relation relation = rltnImdg.getFirstObjectByPredicate(prdct);
if (relation == null || (!ServiceStatus.Active.equalsByKey(relation.getServiceStatus()) && !ServiceStatus.Reopened.equalsByKey(relation.getServiceStatus()))) {
return of(ClearingError.ClearingUnavailableForCompany, validatedObject.getSection(), validatedObject.getCompanyId());
}
context.storeObject(RegistryValidationStored.Relation, relation);
return empty();
}
}

View file

@ -0,0 +1,34 @@
package ru.spcex.clearing.service.validation.admission;
import java.util.Optional;
import ru.clearing.classes.statics.data.company.Company;
import ru.clearing.classes.statics.data.registry.Registry;
import ru.spcex.clearing.error.ClearingError;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.platform.enumeration.WorkflowStatus;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.validation.ImdgValidationContext;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.validation.IValidationRule;
public class CompanyActiveValidationRule implements IValidationRule<ImdgValidationContext<Registry>> {
@Override
public String ruleName() {
return "CompanyActive";
}
@Override
public Optional<EnumMessage> validate(ImdgValidationContext<Registry> context) {
Registry validatedObject = context.getValidatedObject();
Imdg<Company> cmpImdg = context.obtainMap(IMDGDistributedNames.Map_Company, Company.class);
if (validatedObject.getCompanyId() == null) {
return of(ClearingError.CompanyNotActive, validatedObject.getCompanyId());
}
Company company = cmpImdg.getSingleObjectByID(validatedObject.getCompanyId());
if (company == null || !WorkflowStatus.Active.equalsByKey(company.getWorkflowStatus())) {
return of(ClearingError.CompanyNotActive, validatedObject.getCompanyId());
}
return empty();
}
}

View file

@ -0,0 +1,33 @@
package ru.spcex.clearing.service.validation.admission;
import java.util.Optional;
import ru.clearing.classes.statics.data.registry.Registry;
import ru.clearing.classes.statics.data.registry.TradingClearingRegistry;
import ru.spcex.clearing.error.ClearingError;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.platform.enumeration.ServiceStatus;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.validation.ImdgValidationContext;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.validation.IValidationRule;
public class TcrActiveValidationRule implements IValidationRule<ImdgValidationContext<Registry>> {
@Override
public String ruleName() {
return "TcrActive";
}
@Override
public Optional<EnumMessage> validate(ImdgValidationContext<Registry> context) {
Registry validatedObject = context.getValidatedObject();
if (validatedObject.getTradingClearingRegistryId() == null) {
return of(ClearingError.TradingClearingRegistryNotActive, validatedObject.getTradingClearingRegistryId());
}
Imdg<TradingClearingRegistry> tradingClearingRegistryImdg = context.obtainMap(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class);
TradingClearingRegistry tcr = tradingClearingRegistryImdg.getSingleObjectByID(validatedObject.getTradingClearingRegistryId());
if (tcr == null || !(ServiceStatus.Active.equalsByKey(tcr.getStatus()) || ServiceStatus.Reopened.equalsByKey(tcr.getStatus())))
return of(ClearingError.TradingClearingRegistryNotActive, validatedObject.getTradingClearingRegistryId());
return empty();
}
}

View file

@ -4,8 +4,12 @@ import java.time.Instant;
import java.time.LocalDate;
import java.util.ArrayList;
import java.util.Collection;
import java.util.Collections;
import java.util.Comparator;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.Optional;
import java.util.stream.Collectors;
import org.slf4j.Logger;
@ -26,6 +30,7 @@ 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.InclusionToPoolPayload;
import ru.spcex.platform.classes.base.interfaces.ExecutionType;
import ru.spcex.platform.enumeration.RegistryDesignation;
import ru.spcex.platform.enumeration.RegistryInstrumentType;
import ru.spcex.platform.enumeration.RegistryStatus;
@ -127,6 +132,9 @@ public class InclusionObligations implements ISessionStage {
.filter(registry -> registry.getSettlementDate().isEqual(LocalDate.now()))
.collect(Collectors.groupingBy(Registry::getGroupId));
Optional<SessionType> ssnType = obtainSessionType(sessionId);
Map<Long, Registry> rgsToUpdate = new HashMap<>();
List<ExecutionCommon> execsToUpdate = new ArrayList<>();
for (Map.Entry<Long, List<Registry>> entrySet : registryByGroupId.entrySet()) {
log.debug("Processing set of registry with groupId: {}", entrySet.getKey());
String rgsSection = null;
@ -134,25 +142,27 @@ public class InclusionObligations implements ISessionStage {
registry.setRegistryStatus(RegistryStatus.POOL.getKey());
registry.setClearingDate(LocalDate.now());
registry.setSessionId(sessionId);
obtainSessionType(sessionId).ifPresent(st -> registry.setSessionType(st.getKey()));
registryImdg.update(registry);
ssnType.ifPresent(st -> registry.setSessionType(st.getKey()));
rgsToUpdate.put(registry.getId(), registry);
if (TextUtil.isEmpty(rgsSection)) {
rgsSection = registry.getSection();
}
}
updateExecutions(sessionId, entrySet.getKey(), rgsSection);
execsToUpdate.addAll(updateExecutions(sessionId, entrySet.getKey(), rgsSection));
}
registryImdg.putAll(rgsToUpdate);
updateExecsBatchV2(execsToUpdate);
return new StageResult(null, true);
}
@SuppressWarnings("unchecked")
private <T extends ExecutionCommon> void updateExecutions(Long sessionId, Long rgsGroupId, String rgsSection) {
private <T extends ExecutionCommon> Collection<ExecutionCommon> updateExecutions(Long sessionId, Long rgsGroupId, String rgsSection) {
Instant now = Instant.now();
SessionType ssnTpe = obtainSessionType(sessionId).orElse(null);
if (ssnTpe == null || sessionId == null) {
log.debug("will not update executions#sessionId - couldn't determine session type");
return;
return Collections.emptyList();
}
Imdg<T> execImdg = null;
switch (ssnTpe) {
@ -173,17 +183,18 @@ public class InclusionObligations implements ISessionStage {
if (execImdg == null) {
log.warn("couldn't define Execution Type for session {}. Will not update executions#sessionId",
sessionType);
return;
return Collections.emptyList();
}
ImdgPredicateBuilder pb = execImdg.predicateBuilder();
ImdgPredicate prdct = pb.equals("exchangeExecutionId", rgsGroupId);
Collection<T> execs = execImdg.getCollectionObjectsByPredicate(prdct);
Imdg<T> finalExecImdg = execImdg;
//Imdg<T> finalExecImdg = execImdg;
execs.forEach(e -> {
e.setUpdated(now);
e.setSessionId(sessionId);
finalExecImdg.update(e);
// finalExecImdg.update(e);
});
return (Collection<ExecutionCommon>) execs;
}
private Optional<SessionType> obtainSessionType(Long sessionId) {
@ -217,4 +228,54 @@ public class InclusionObligations implements ISessionStage {
return rgsPb.alwaysTrue();
}
}
private static <T extends ExecutionCommon> Map<Long, T> castMap(Map<Long, ExecutionCommon> m) {
return (Map<Long, T>) m;
// ExecutionType type = m.values().stream().map(e -> e.type()).findFirst().orElse(null);
// if (type == null) {
// return Collections.emptyMap();
// }
// switch (type) {
// case ExecutionDeposit -> {
// return (Map<Long, T>) m;
// }
// case ExecutionCurrency -> {
// return (Map<Long, T>) m;
// }
// }
}
private void updateExecsBatchV2(List<ExecutionCommon> execs) {
log.debug("updating {} executions", execs.size());
execs.sort(Comparator.comparing(ExecutionCommon::type));
log.debug("sorted executions by type");
Map<Long, ExecutionCommon> m = new HashMap<>();
for (int i = 0; i < execs.size(); i++) {
ExecutionCommon exec = execs.get(i);
ExecutionType execType = exec.type();
log.debug("batch from {} position type {}", i, execType);
m.put(exec.getId(), exec);
int j = i + 1;
while (j < execs.size() && j < i + 100) {
ExecutionCommon other = execs.get(j);
if (!Objects.equals(other.type(), execType)) {
i = j;
break;
}
m.put(other.getId(), other);
j++;
}
log.debug("batch size {} type {}", m.size(), execType);
if (!m.isEmpty()) {
switch (execType) {
case ExecutionDeposit -> executionDepositImdg.putAll(castMap(m));
case ExecutionFond -> executionFondImdg.putAll(castMap(m));
case ExecutionCurrency -> executionCurrImdg.putAll(castMap(m));
}
}
log.debug("batch size {} type {} done", m.size(), execType);
m.clear();
}
}
}

View file

@ -1,14 +1,28 @@
package ru.spcex.clearing.session.stage.impl;
import java.util.ArrayList;
import java.util.Collection;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.stream.Collectors;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.beans.factory.config.ConfigurableBeanFactory;
import org.springframework.context.annotation.Scope;
import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.account.Account;
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.imdg.IMDGDistributedNames;
import ru.spcex.clearing.service.validation.admission.AccountActiveValidationRule;
import ru.spcex.clearing.service.validation.admission.ClearingAvailableValidationRule;
import ru.spcex.clearing.service.validation.admission.CompanyActiveValidationRule;
import ru.spcex.clearing.service.validation.admission.TcrActiveValidationRule;
import ru.spcex.clearing.session.stage.ISessionStage;
import ru.spcex.clearing.session.stage.StageResult;
import ru.spcex.clearing.session.stage.Task;
@ -16,41 +30,46 @@ import ru.spcex.clearing.session.stage.TaskType;
import ru.spcex.platform.enumeration.RegistryStatus;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.imdg.iml.hazelcast.adapter.CashCloser;
import ru.spcex.platform.imdg.iml.hazelcast.adapter.CashV2ById;
import ru.spcex.platform.imdg.iml.hazelcast.adapter.CashV2ByIdAndString;
import ru.spcex.platform.imdg.validation.ImdgValidationContext;
import ru.spcex.platform.imdg.validation.LogPrefixId;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.enumeration.IMessageResolver;
import ru.spcex.platform.utils.validation.IValidator;
import java.util.Collection;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.function.Function;
import java.util.stream.Collectors;
import ru.spcex.platform.utils.validation.ValidatorImpl;
@Service
@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE)
public class ObligationAdmission implements ISessionStage {
private final Logger log = LoggerFactory.getLogger(getClass());
private final Imdg<Registry> registryImdg;
private final ImdgProvider imdgProvider;
private final Function<Registry, IValidator> validatorFactory;
private final Imdg<Relation> imdgRltn;
private final Imdg<Account> imdgAcc;
private final Imdg<Company> imdgCmp;
private final Imdg<TradingClearingRegistry> imdgTcr;
private final IMessageResolver messageResolver;
private final ObligationAdmissionCash cash = new ObligationAdmissionCash();
@Autowired
public ObligationAdmission(ImdgProvider imdgProvider,
@Qualifier("obligationAndRequirementsAdmissionValidator") Function<Registry, IValidator> validatorFactory,
IMessageResolver messageResolver) {
this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
this.imdgProvider = imdgProvider;
this.validatorFactory = validatorFactory;
this.messageResolver = messageResolver;
this.imdgRltn = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_Relation, Relation.class, null);
this.imdgAcc = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_Account, Account.class, cash.accCash);
this.imdgCmp = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_Company, Company.class, cash.cmpCash);
this.imdgTcr = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class, cash.tcrCash);
}
@SuppressWarnings("unchecked")
@Override
public StageResult<?> submit(Task<?> task) {
if (task.getTaskType() == TaskType.ObligationsAdmission) {
return obligationAdmission((Long) task.getData());
try (cash) {
return obligationAdmission((Long) task.getData());
}
}
throw new IllegalStateException("Unknown task type: " + task.getTaskType());
}
@ -59,11 +78,12 @@ public class ObligationAdmission implements ISessionStage {
Collection<Registry> registries = registryImdg.getCollectionObjectsBySQL("sessionId = " + sessionId);
Map<Long, List<Registry>> byGroups = registries.stream().collect(Collectors.groupingBy(Registry::getGroupId));
log.info("found {} ({} groups) registries with sessionId {}", registries.size(), byGroups.size(), sessionId);
Map<Long, Registry> rgsToUpdate = new HashMap<>();
for (Map.Entry<Long, List<Registry>> grpEntry : byGroups.entrySet()) {
EnumMessage groupError = null;
List<Registry> rgsGroup = grpEntry.getValue();
for (Registry rgs : rgsGroup) {
IValidator validator = validatorFactory.apply(rgs);
IValidator validator = valFor(rgs);
Optional<EnumMessage> error = validator.tillFirstError();
if (error.isPresent()) {
groupError = error.get();
@ -74,10 +94,44 @@ public class ObligationAdmission implements ISessionStage {
log.warn("error {} for registries groupId = {}", messageResolver.resolve(groupError), grpEntry.getKey());
for (Registry rgs : rgsGroup) {
rgs.setRegistryStatus(RegistryStatus.NACK.getKey());
registryImdg.update(rgs);
rgsToUpdate.put(rgs.getId(), rgs);
}
}
}
registryImdg.putAll(rgsToUpdate);
return new StageResult<>(null, true);
}
private IValidator valFor(Registry rgs) {
ImdgValidationContext<Registry> ctx = new ImdgValidationContext<>();
ctx.setValidatedObject(rgs);
ctx.addImdg(IMDGDistributedNames.Map_Relation, imdgRltn);
//ctx.addImdg(IMDGDistributedNames.Map_Session, imdgSession);
//ctx.addImdg(IMDGDistributedNames.Map_SectionDictionary, imdgSectionDictionary);
ctx.addImdg(IMDGDistributedNames.Map_Account, imdgAcc);
ctx.addImdg(IMDGDistributedNames.Map_Company, imdgCmp);
ctx.addImdg(IMDGDistributedNames.Map_TradingClearingRegistry, imdgTcr);
ctx.setLogPrefix(LogPrefixId.INSTANCE);
return new ValidatorImpl<>(ctx,
new ClearingAvailableValidationRule(cash.rltnCash),
new AccountActiveValidationRule(),
new CompanyActiveValidationRule(),
new TcrActiveValidationRule());
}
protected static class ObligationAdmissionCash extends CashCloser {
CashV2ByIdAndString<Relation> rltnCash;
CashV2ById<Account> accCash;
CashV2ById<Company> cmpCash;
CashV2ById<TradingClearingRegistry> tcrCash;
public ObligationAdmissionCash() {
this.cashes = new ArrayList<>();
rltnCash = add(new CashV2ByIdAndString<>("rltnCash", rltn ->
new CashV2ByIdAndString.CustomKey(rltn.getConsumerId(), rltn.getService())));
accCash = add(new CashV2ById<>("accCash"));
cmpCash = add(new CashV2ById<>("cmpCash"));
tcrCash = add(new CashV2ById<>("tcrCash"));
}
}
}

View file

@ -22,6 +22,10 @@ public class ExecutionCurrencyMapStore extends TemplateMapStore<ExecutionCurrenc
return IMDGDistributedNames.Map_ExecutionCurrency;
}
public String[] getIndexingField() {
return new String[]{"exchangeExecutionId", "sessionId"};
}
@Override
public String getTableName() {
return "EXECUTION_CURRENCY";