From 46da00357ec002aa3f0c1432c607b3b12dc515d9 Mon Sep 17 00:00:00 2001 From: ialbert Date: Wed, 31 Jul 2024 13:56:34 +0300 Subject: [PATCH] sessions step 3/4 optimization imdg execution currency index --- ...yConsumerIdAndServiceCashingPredicate.java | 30 +++++ .../clearing/config/ValidationConfig.java | 20 ---- .../RegistryStep3ValidationRule.java | 107 ------------------ .../AccountActiveValidationRule.java | 32 ++++++ .../ClearingAvailableValidationRule.java | 70 ++++++++++++ .../CompanyActiveValidationRule.java | 34 ++++++ .../admission/TcrActiveValidationRule.java | 33 ++++++ .../stage/impl/InclusionObligations.java | 77 +++++++++++-- .../stage/impl/ObligationAdmission.java | 86 +++++++++++--- .../ExecutionCurrencyMapStore.java | 4 + 10 files changed, 342 insertions(+), 151 deletions(-) create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/component/predicate/cash/relation/RelationByConsumerIdAndServiceCashingPredicate.java delete mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/RegistryStep3ValidationRule.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/admission/AccountActiveValidationRule.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/admission/ClearingAvailableValidationRule.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/admission/CompanyActiveValidationRule.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/admission/TcrActiveValidationRule.java diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/component/predicate/cash/relation/RelationByConsumerIdAndServiceCashingPredicate.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/component/predicate/cash/relation/RelationByConsumerIdAndServiceCashingPredicate.java new file mode 100644 index 000000000..de1708aad --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/component/predicate/cash/relation/RelationByConsumerIdAndServiceCashingPredicate.java @@ -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 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; + } + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/ValidationConfig.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/ValidationConfig.java index 9271670e1..c872a5bd5 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/ValidationConfig.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/ValidationConfig.java @@ -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 registryValidator() { - return rgs -> { - ImdgValidationContext 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 sdf01Validator() { return sDf01 -> { diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/RegistryStep3ValidationRule.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/RegistryStep3ValidationRule.java deleted file mode 100644 index c60eb806a..000000000 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/RegistryStep3ValidationRule.java +++ /dev/null @@ -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> { - ClearingAvailable() { - @Override - public Optional validate(ImdgValidationContext context) { - Registry validatedObject = context.getValidatedObject(); -// Supplier sectionFind = () -> { -// Imdg sessionImdg = context.obtainMap(IMDGDistributedNames.Map_Session, Session.class); -// Imdg 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 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 validate(ImdgValidationContext context) { - Registry validatedObject = context.getValidatedObject(); - Imdg 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 validate(ImdgValidationContext context) { - Registry validatedObject = context.getValidatedObject(); - Imdg 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 validate(ImdgValidationContext context) { - Registry validatedObject = context.getValidatedObject(); - if (validatedObject.getTradingClearingRegistryId() == null) { - return of(ClearingError.TradingClearingRegistryNotActive, validatedObject.getTradingClearingRegistryId()); - } - Imdg 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(); - } -} \ No newline at end of file diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/admission/AccountActiveValidationRule.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/admission/AccountActiveValidationRule.java new file mode 100644 index 000000000..fc39d34a9 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/admission/AccountActiveValidationRule.java @@ -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> { + + @Override + public String ruleName() { + return "AccountActive"; + } + + @Override + public Optional validate(ImdgValidationContext context) { + Registry validatedObject = context.getValidatedObject(); + Imdg 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(); + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/admission/ClearingAvailableValidationRule.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/admission/ClearingAvailableValidationRule.java new file mode 100644 index 000000000..f3adbee77 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/admission/ClearingAvailableValidationRule.java @@ -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> { + private final CashV2ByIdAndString rltnCash; + + public ClearingAvailableValidationRule(CashV2ByIdAndString rltnCash) { + this.rltnCash = rltnCash; + } + + @Override + public String ruleName() { + return "ClearingAvailable"; + } + + @Override + public Optional validate(ImdgValidationContext context) { + Registry validatedObject = context.getValidatedObject(); + //Supplier sectionFind = () -> { + // Imdg sessionImdg = context.obtainMap(IMDGDistributedNames.Map_Session, Session.class); + // Imdg 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 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(); + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/admission/CompanyActiveValidationRule.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/admission/CompanyActiveValidationRule.java new file mode 100644 index 000000000..14b83ff2e --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/admission/CompanyActiveValidationRule.java @@ -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> { + + @Override + public String ruleName() { + return "CompanyActive"; + } + + @Override + public Optional validate(ImdgValidationContext context) { + Registry validatedObject = context.getValidatedObject(); + Imdg 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(); + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/admission/TcrActiveValidationRule.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/admission/TcrActiveValidationRule.java new file mode 100644 index 000000000..07d7c72db --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/admission/TcrActiveValidationRule.java @@ -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> { + + @Override + public String ruleName() { + return "TcrActive"; + } + + @Override + public Optional validate(ImdgValidationContext context) { + Registry validatedObject = context.getValidatedObject(); + if (validatedObject.getTradingClearingRegistryId() == null) { + return of(ClearingError.TradingClearingRegistryNotActive, validatedObject.getTradingClearingRegistryId()); + } + Imdg 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(); + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InclusionObligations.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InclusionObligations.java index 877044b98..c1b1b1706 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InclusionObligations.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InclusionObligations.java @@ -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 ssnType = obtainSessionType(sessionId); + Map rgsToUpdate = new HashMap<>(); + List execsToUpdate = new ArrayList<>(); for (Map.Entry> 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 void updateExecutions(Long sessionId, Long rgsGroupId, String rgsSection) { + private Collection 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 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 execs = execImdg.getCollectionObjectsByPredicate(prdct); - Imdg finalExecImdg = execImdg; + //Imdg finalExecImdg = execImdg; execs.forEach(e -> { e.setUpdated(now); e.setSessionId(sessionId); - finalExecImdg.update(e); +// finalExecImdg.update(e); }); + return (Collection) execs; } private Optional obtainSessionType(Long sessionId) { @@ -217,4 +228,54 @@ public class InclusionObligations implements ISessionStage { return rgsPb.alwaysTrue(); } } + + private static Map castMap(Map m) { + return (Map) 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) m; +// } +// case ExecutionCurrency -> { +// return (Map) m; +// } +// } + } + + private void updateExecsBatchV2(List execs) { + log.debug("updating {} executions", execs.size()); + execs.sort(Comparator.comparing(ExecutionCommon::type)); + log.debug("sorted executions by type"); + + Map 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(); + } + } } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/ObligationAdmission.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/ObligationAdmission.java index b8f911b15..a422134d6 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/ObligationAdmission.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/ObligationAdmission.java @@ -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 registryImdg; - private final ImdgProvider imdgProvider; - private final Function validatorFactory; + private final Imdg imdgRltn; + private final Imdg imdgAcc; + private final Imdg imdgCmp; + private final Imdg imdgTcr; private final IMessageResolver messageResolver; + private final ObligationAdmissionCash cash = new ObligationAdmissionCash(); @Autowired public ObligationAdmission(ImdgProvider imdgProvider, - @Qualifier("obligationAndRequirementsAdmissionValidator") Function 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 registries = registryImdg.getCollectionObjectsBySQL("sessionId = " + sessionId); Map> byGroups = registries.stream().collect(Collectors.groupingBy(Registry::getGroupId)); log.info("found {} ({} groups) registries with sessionId {}", registries.size(), byGroups.size(), sessionId); + Map rgsToUpdate = new HashMap<>(); for (Map.Entry> grpEntry : byGroups.entrySet()) { EnumMessage groupError = null; List rgsGroup = grpEntry.getValue(); for (Registry rgs : rgsGroup) { - IValidator validator = validatorFactory.apply(rgs); + IValidator validator = valFor(rgs); Optional 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 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 rltnCash; + CashV2ById accCash; + CashV2ById cmpCash; + CashV2ById 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")); + } + } } diff --git a/clearing-parent/imdg/src/main/java/ru/spcex/clearing/imdg/businessobject/ExecutionCurrencyMapStore.java b/clearing-parent/imdg/src/main/java/ru/spcex/clearing/imdg/businessobject/ExecutionCurrencyMapStore.java index 7185976ea..3b4a76ed8 100644 --- a/clearing-parent/imdg/src/main/java/ru/spcex/clearing/imdg/businessobject/ExecutionCurrencyMapStore.java +++ b/clearing-parent/imdg/src/main/java/ru/spcex/clearing/imdg/businessobject/ExecutionCurrencyMapStore.java @@ -22,6 +22,10 @@ public class ExecutionCurrencyMapStore extends TemplateMapStore