From 238cbf4f8e3f59964aac34234eb61bc107e777f1 Mon Sep 17 00:00:00 2001 From: ialbert Date: Thu, 8 Aug 2024 16:15:29 +0300 Subject: [PATCH] execution load optimization cash --- .../clearing/config/ValidationConfig.java | 61 ----- .../execution/ExecutionCurrencyComponent.java | 250 +++++++++++------- .../execution/ExecutionDepositComponent.java | 228 +++++++++------- .../execution/ExecutionFondComponent.java | 248 ++++++++++------- ...anyByTradingCodeCashingValidationRule.java | 66 +++++ ...sSecurityPresentDepositValidationRule.java | 52 ++++ .../STradesSecurityPresentValidationRule.java | 57 ++++ ...deAndCmpIsActiveCashingValidationRule.java | 71 +++++ .../impl/FormingPaymentInstructionAssets.java | 1 + .../predicate/specific/SecuritySelector.java | 42 ++- 10 files changed, 740 insertions(+), 336 deletions(-) create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/strades/CompanyByTradingCodeCashingValidationRule.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/strades/STradesSecurityPresentDepositValidationRule.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/strades/STradesSecurityPresentValidationRule.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/strades/TcrByCodeAndCmpIsActiveCashingValidationRule.java 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 c872a5bd5..883481a8d 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 @@ -16,7 +16,6 @@ import ru.clearing.classes.statics.data.execution.ExecutionDeposit; import ru.clearing.classes.statics.data.instrument.issue.EquitySecurity; import ru.clearing.classes.statics.data.instrument.issue.FixedIncomeSecurity; import ru.clearing.classes.statics.data.misc.Currency; -import ru.clearing.classes.statics.data.misc.STrades; import ru.clearing.classes.statics.data.misc.Session; import ru.clearing.classes.statics.data.payment.PaymentInstruction; import ru.clearing.classes.statics.data.registry.Registry; @@ -55,7 +54,6 @@ 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.ReturnDepositValidationRule; -import ru.spcex.clearing.service.validation.STradesValidationRule; import ru.spcex.clearing.service.validation.Sdf01NewValidationRule; import ru.spcex.clearing.service.validation.Sdf01ValidationRule; import ru.spcex.clearing.service.validation.Sdf06NewValidationRule; @@ -169,65 +167,6 @@ public class ValidationConfig { }; } - @Bean("sTradesValidatorFond") - public Function sTradesValidator() { - return sTrades -> { - ImdgValidationContext context = new ImdgValidationContext<>(); - context.setValidatedObject(sTrades); - context.addImdg(IMDGDistributedNames.Map_Security, imdgSecurity); - context.addImdg(IMDGDistributedNames.Map_Company, imdgCompany); - context.addImdg(IMDGDistributedNames.Map_TradingClearingRegistry, imdgTradingClearingRegistry); - context.addImdg(IMDGDistributedNames.Map_FixedIncomeSecurity, fixedIncomeSecurityImdg); - context.addImdg(IMDGDistributedNames.Map_MoneyMarketSecurity, imdgMoneyMarketSecurity); - context.addImdg(IMDGDistributedNames.Map_EquitySecurity, equitySecurityImdg); - context.addImdg(IMDGDistributedNames.Map_CurrencyPairSecurity, currencyPairSecurityImdg); - context.setLogPrefix(LogPrefixId.INSTANCE); - return new ValidatorImpl<>(context, - STradesValidationRule.SecurityPresentFond, - STradesValidationRule.CompanyPresent, - STradesValidationRule.CounterCompanyPresent, - STradesValidationRule.TradingClearingRegistryPresent); - }; - } - - @Bean("sTradesValidatorDeposit") - public Function sTradesValidatorDeposit() { - return sTrades -> { - ImdgValidationContext context = new ImdgValidationContext<>(); - context.setValidatedObject(sTrades); - context.addImdg(IMDGDistributedNames.Map_Security, imdgSecurity); - context.addImdg(IMDGDistributedNames.Map_MoneyMarketSecurity, imdgMoneyMarketSecurity); - context.addImdg(IMDGDistributedNames.Map_Company, imdgCompany); - context.addImdg(IMDGDistributedNames.Map_TradingClearingRegistry, imdgTradingClearingRegistry); - context.setLogPrefix(LogPrefixId.INSTANCE); - return new ValidatorImpl<>(context, - STradesValidationRule.SecurityPresentDeposit, - STradesValidationRule.CompanyPresent, - STradesValidationRule.CounterCompanyPresent, - STradesValidationRule.TradingClearingRegistryPresent); - }; - } - - @Bean("sTradesValidatorCurrency") - public Function sTradesValidatorCurrency() { - return sTrades -> { - ImdgValidationContext context = new ImdgValidationContext<>(); - context.setValidatedObject(sTrades); - context.addImdg(IMDGDistributedNames.Map_CurrencyPairSecurity, currencyPairSecurityImdg); - context.addImdg(IMDGDistributedNames.Map_MoneyMarketSecurity, imdgMoneyMarketSecurity); - context.addImdg(IMDGDistributedNames.Map_FixedIncomeSecurity, fixedIncomeSecurityImdg); - context.addImdg(IMDGDistributedNames.Map_EquitySecurity, equitySecurityImdg); - context.addImdg(IMDGDistributedNames.Map_Company, imdgCompany); - context.addImdg(IMDGDistributedNames.Map_TradingClearingRegistry, imdgTradingClearingRegistry); - context.setLogPrefix(LogPrefixId.INSTANCE); - return new ValidatorImpl<>(context, - STradesValidationRule.SecurityPresentCurrencyPair, - STradesValidationRule.CompanyPresent, - STradesValidationRule.CounterCompanyPresent, - STradesValidationRule.TradingClearingRegistryPresent); - }; - } - @Bean("sdf01Validator") public Function sdf01Validator() { return sDf01 -> { diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/execution/ExecutionCurrencyComponent.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/execution/ExecutionCurrencyComponent.java index 8d3eec2f4..72ac45058 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/execution/ExecutionCurrencyComponent.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/execution/ExecutionCurrencyComponent.java @@ -8,7 +8,6 @@ import java.util.Collection; import java.util.HashMap; import java.util.Map; import java.util.Optional; -import java.util.function.Function; import org.apache.kafka.clients.producer.Producer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -18,9 +17,13 @@ import org.springframework.scheduling.annotation.EnableScheduling; import org.springframework.stereotype.Component; import ru.clearing.classes.statics.data.company.Company; import ru.clearing.classes.statics.data.execution.ExecutionCurrency; +import ru.clearing.classes.statics.data.instrument.issue.EquitySecurity; +import ru.clearing.classes.statics.data.instrument.issue.FixedIncomeSecurity; import ru.clearing.classes.statics.data.misc.Listing; import ru.clearing.classes.statics.data.misc.STrades; import ru.clearing.classes.statics.data.registry.TradingClearingRegistry; +import ru.clearing.classes.statics.data.security.CurrencyPairSecurity; +import ru.clearing.classes.statics.data.security.MoneyMarketSecurity; import ru.clearing.classes.statics.data.security.Security; import ru.spcex.clearing.error.ClearingError; import ru.spcex.clearing.error.ClearingException; @@ -30,6 +33,9 @@ import ru.spcex.clearing.platform.messaging.domain.cud.registry.DealRegisterNewR import ru.spcex.clearing.platform.messaging.domain.cud.registry.ExecutionType; import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.clearing.service.validation.ValidationStored; +import ru.spcex.clearing.service.validation.strades.CompanyByTradingCodeCashingValidationRule; +import ru.spcex.clearing.service.validation.strades.STradesSecurityPresentValidationRule; +import ru.spcex.clearing.service.validation.strades.TcrByCodeAndCmpIsActiveCashingValidationRule; import ru.spcex.platform.enumeration.Section; import ru.spcex.platform.enumeration.Side; import ru.spcex.platform.imdg.api.Imdg; @@ -37,6 +43,11 @@ import ru.spcex.platform.imdg.api.ImdgId; 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.iml.hazelcast.adapter.CashCloser; +import ru.spcex.platform.imdg.iml.hazelcast.adapter.CashV2ByIdAndString; +import ru.spcex.platform.imdg.iml.hazelcast.adapter.CashV2ByString; +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.IEnumId; import ru.spcex.platform.utils.enumeration.IEnumKey; @@ -44,6 +55,7 @@ import ru.spcex.platform.utils.enumeration.IMessageResolver; import ru.spcex.platform.utils.log.ExceptionUtils; import ru.spcex.platform.utils.time.TimeUtil; import ru.spcex.platform.utils.validation.IValidator; +import ru.spcex.platform.utils.validation.ValidatorImpl; /** * ExecutionCurrency - Сделки на Валютной секции @@ -55,13 +67,19 @@ public class ExecutionCurrencyComponent implements IExecutionUploadComponent { private final Imdg sTradeImdg; private final Imdg executionCurrencyImdg; + private final Imdg currencyPairSecurityImdg; + private final Imdg imdgMoneyMarketSecurity; + private final Imdg fixedIncomeSecurityImdg; + private final Imdg equitySecurityImdg; + private final Imdg cmpImdg; + private final Imdg tcrImdg; + private final ExecCurrCompCash valCash = new ExecCurrCompCash(); private final ImdgId idGen; private final Imdg listingImdg; // protected transient Long tradeNum; //todo used? // protected transient Instant tradingDay; private final IMessageResolver msgResolver; - private final Function stradesValidator; private final KafkaSender kafkaSender; private static final DateTimeFormatter contractFormatter = DateTimeFormatter.ofPattern("ddMMyy"); private final Map cash = new HashMap<>(); @@ -70,15 +88,19 @@ public class ExecutionCurrencyComponent implements IExecutionUploadComponent { @Autowired public ExecutionCurrencyComponent(ImdgProvider imdgProvider, Producer kafka, - @Qualifier("sTradesValidatorCurrency") Function stradesValidator, @Qualifier("kafkaSenderWithoutRequestInfo") KafkaSender kafkaSender, IMessageResolver msgResolver) { this.sTradeImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_STrades, STrades.class); this.executionCurrencyImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionCurrency, ExecutionCurrency.class); this.listingImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Listing, Listing.class); - this.stradesValidator = stradesValidator; this.kafkaSender = kafkaSender; this.msgResolver = msgResolver; + this.currencyPairSecurityImdg = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_CurrencyPairSecurity, CurrencyPairSecurity.class, null); + this.imdgMoneyMarketSecurity = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_MoneyMarketSecurity, MoneyMarketSecurity.class, null); + this.fixedIncomeSecurityImdg = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_FixedIncomeSecurity, FixedIncomeSecurity.class, null); + this.equitySecurityImdg = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_EquitySecurity, EquitySecurity.class, null); + this.cmpImdg = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_Company, Company.class, null); + this.tcrImdg = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class, null); // resetTradingDay(); this.idGen = imdgProvider.getImdgIdGenerator(); @@ -104,95 +126,97 @@ public class ExecutionCurrencyComponent implements IExecutionUploadComponent { // } public void processNewTS(Long fromId) { - boolean cleanLoad = fromId == null; - if (cleanLoad) { - log.info("clearing cash..."); - cash.clear(); - initExecCash(); - } - LocalDate today = LocalDate.now(); - log.debug("Start check new S_TRADE at {}, fromId={}", today, fromId); - //выбираем STrades на сегодня с правильным section - Collection sTrades; - if (cleanLoad) { - sTrades = sTradeImdg.getCollectionObjectsByFieldValues(Map.of( - "tradeDate", LocalDate.now(), - "section", Section.CURR.getKey() - )); - } else { - ImdgPredicateBuilder pb = sTradeImdg.predicateBuilder(); - ImdgPredicate filter = pb.and( - pb.equals("tradeDate", LocalDate.now()), - pb.equals("section", Section.CURR.getKey()), - pb.greatEqual("id", fromId) - ); - sTrades = sTradeImdg.getCollectionObjectsByPredicate(filter); - } - log.info("Found {} s_trade for today (fromId={})", sTrades.size(), fromId); - - if (sTrades.isEmpty()) { - logError(ClearingError.NewDealsNotFound); - return; - } - - //убираем уже добавленные в ExecutionCurrency - if (!cleanLoad) { - log.info("filtering strades based on ExecutionCurrency cash..."); - sTrades.removeIf(sTrd -> { - Side sTrdSide = IEnumKey.getEnumByKey(Side.class, sTrd.getOperation()); - if (sTrdSide == null || !(sTrdSide.equals(Side.BUY) || sTrdSide.equals(Side.SELL))) { - log.warn("cannot parse strade.tradeNum={} strade.operation {}", sTrd.getTradeNum(), sTrd.getOperation()); - return true; - } -// MoneyFlowSide excDepSide = sTrdSide == Side.BUY ? MoneyFlowSide.BUY : MoneyFlowSide.SELL; - ExecUploadKey excKey = ExecUploadKey.cash(sTrd.getTradeDate(), - sTrd.getTradeNum(), - sTrdSide.getKey()); - return cash.containsKey(excKey); - }); - } - log.info("{} strades left after already-added filtering", sTrades.size()); - - Map execsToInsert = new HashMap<>(); - for (STrades sTrd : sTrades) { - log.trace("S_TRADE[{}] new", sTrd.getId()); - - //проверки - IValidator validator = stradesValidator.apply(sTrd); - Optional error = validator.tillFirstError(); - if (error.isPresent()) { - logError(sTrd.getId(), error.get()); - continue; + try (valCash) { + boolean cleanLoad = fromId == null; + if (cleanLoad) { + log.info("clearing cash..."); + cash.clear(); + initExecCash(); } - ExecutionCurrency newEC; - try { - newEC = createExecutionCurrency(sTrd, validator); - ExecUploadKey execKey = ExecUploadKey.cash(newEC); - ExecUploadCashInfo storedId; - if (!cleanLoad || (storedId = cash.get(execKey)) == null) { - newEC.setId(idGen.nextId()); - cash.put(execKey, ExecUploadCashInfo.cash(newEC)); - log.debug("New executionCurrency.id={} was created.", newEC.getId()); - } else { - newEC.setId(storedId.id()); - newEC.setUpdated(newEC.getCreated()); - newEC.setCreated(storedId.createDt()); - log.debug("executionCurrency.id={} was updated.", newEC.getId()); - } - execsToInsert.put(newEC.getId(), newEC); - //executionCurrencyImdg.insert(newEC); - sendNotification(newEC); - } catch (ClearingException ce) { - auditMessage(ce); - } catch (Exception e) { - log.error("When create new ExecutionCurrency by STrade[{}] error: {}", sTrd.getId(), ExceptionUtils.getStackTrace(e)); + LocalDate today = LocalDate.now(); + log.debug("Start check new S_TRADE at {}, fromId={}", today, fromId); + //выбираем STrades на сегодня с правильным section + Collection sTrades; + if (cleanLoad) { + sTrades = sTradeImdg.getCollectionObjectsByFieldValues(Map.of( + "tradeDate", LocalDate.now(), + "section", Section.CURR.getKey() + )); + } else { + ImdgPredicateBuilder pb = sTradeImdg.predicateBuilder(); + ImdgPredicate filter = pb.and( + pb.equals("tradeDate", LocalDate.now()), + pb.equals("section", Section.CURR.getKey()), + pb.greatEqual("id", fromId) + ); + sTrades = sTradeImdg.getCollectionObjectsByPredicate(filter); } - } - log.info("batch putAll {} ExecutionCurrency", execsToInsert.size()); - executionCurrencyImdg.putAll(execsToInsert, 200); + log.info("Found {} s_trade for today (fromId={})", sTrades.size(), fromId); - Long newMaxTradeNum = sTrades.stream().mapToLong(STrades::getTradeNum).max().orElse(0); // orElseGet(() -> tradeNum) - log.info("Process completed. Next tradeNum is {}", newMaxTradeNum); + if (sTrades.isEmpty()) { + logError(ClearingError.NewDealsNotFound); + return; + } + + //убираем уже добавленные в ExecutionCurrency + if (!cleanLoad) { + log.info("filtering strades based on ExecutionCurrency cash..."); + sTrades.removeIf(sTrd -> { + Side sTrdSide = IEnumKey.getEnumByKey(Side.class, sTrd.getOperation()); + if (sTrdSide == null || !(sTrdSide.equals(Side.BUY) || sTrdSide.equals(Side.SELL))) { + log.warn("cannot parse strade.tradeNum={} strade.operation {}", sTrd.getTradeNum(), sTrd.getOperation()); + return true; + } + // MoneyFlowSide excDepSide = sTrdSide == Side.BUY ? MoneyFlowSide.BUY : MoneyFlowSide.SELL; + ExecUploadKey excKey = ExecUploadKey.cash(sTrd.getTradeDate(), + sTrd.getTradeNum(), + sTrdSide.getKey()); + return cash.containsKey(excKey); + }); + } + log.info("{} strades left after already-added filtering", sTrades.size()); + + Map execsToInsert = new HashMap<>(); + for (STrades sTrd : sTrades) { + log.trace("S_TRADE[{}] new", sTrd.getId()); + + //проверки + IValidator validator = valFor(sTrd); + Optional error = validator.tillFirstError(); + if (error.isPresent()) { + logError(sTrd.getId(), error.get()); + continue; + } + ExecutionCurrency newEC; + try { + newEC = createExecutionCurrency(sTrd, validator); + ExecUploadKey execKey = ExecUploadKey.cash(newEC); + ExecUploadCashInfo storedId; + if (!cleanLoad || (storedId = cash.get(execKey)) == null) { + newEC.setId(idGen.nextId()); + cash.put(execKey, ExecUploadCashInfo.cash(newEC)); + log.debug("New executionCurrency.id={} was created.", newEC.getId()); + } else { + newEC.setId(storedId.id()); + newEC.setUpdated(newEC.getCreated()); + newEC.setCreated(storedId.createDt()); + log.debug("executionCurrency.id={} was updated.", newEC.getId()); + } + execsToInsert.put(newEC.getId(), newEC); + //executionCurrencyImdg.insert(newEC); + sendNotification(newEC); + } catch (ClearingException ce) { + auditMessage(ce); + } catch (Exception e) { + log.error("When create new ExecutionCurrency by STrade[{}] error: {}", sTrd.getId(), ExceptionUtils.getStackTrace(e)); + } + } + log.info("batch putAll {} ExecutionCurrency", execsToInsert.size()); + executionCurrencyImdg.putAll(execsToInsert, 200); + + Long newMaxTradeNum = sTrades.stream().mapToLong(STrades::getTradeNum).max().orElse(0); // orElseGet(() -> tradeNum) + log.info("Process completed. Next tradeNum is {}", newMaxTradeNum); + } } @@ -319,7 +343,7 @@ public class ExecutionCurrencyComponent implements IExecutionUploadComponent { } private Optional searchTcrByStrades(STrades sTrades) { - IValidator counterValidator = stradesValidator.apply(sTrades); + IValidator counterValidator = valFor(sTrades); Optional err = counterValidator.tillFirstError(); if (err.isPresent()) { log.error("couldn't extract setCounterPartyTradingClearingRegistryId from strades.id={}", sTrades.getId()); @@ -327,4 +351,48 @@ public class ExecutionCurrencyComponent implements IExecutionUploadComponent { } return Optional.ofNullable(counterValidator.getStored(ValidationStored.STradesTradingClearingRegistry)); } + + public IValidator valFor(STrades sTrades) { + ImdgValidationContext context = new ImdgValidationContext<>(); + context.setValidatedObject(sTrades); + context.addImdg(IMDGDistributedNames.Map_CurrencyPairSecurity, currencyPairSecurityImdg); + context.addImdg(IMDGDistributedNames.Map_MoneyMarketSecurity, imdgMoneyMarketSecurity); + context.addImdg(IMDGDistributedNames.Map_FixedIncomeSecurity, fixedIncomeSecurityImdg); + context.addImdg(IMDGDistributedNames.Map_EquitySecurity, equitySecurityImdg); + context.addImdg(IMDGDistributedNames.Map_Company, cmpImdg); + context.addImdg(IMDGDistributedNames.Map_TradingClearingRegistry, tcrImdg); + context.setLogPrefix(LogPrefixId.INSTANCE); + return new ValidatorImpl<>(context, + new STradesSecurityPresentValidationRule(valCash.fixIncSecCash, + valCash.mmsCash, + valCash.eqtySecCash, + valCash.currPairSecCash), + new CompanyByTradingCodeCashingValidationRule<>( + STrades::getFirmId, valCash.cmpCash, ValidationStored.STradesCompany + ), + new CompanyByTradingCodeCashingValidationRule<>( + STrades::getCpFirmId, valCash.cmpCash, ValidationStored.STradesCompany + ), + new TcrByCodeAndCmpIsActiveCashingValidationRule(valCash.tcrCash) + ); + } + + + private static class ExecCurrCompCash extends CashCloser { + private final CashV2ByString fixIncSecCash; + private final CashV2ByString mmsCash; + private final CashV2ByString eqtySecCash; + private final CashV2ByString currPairSecCash; + private final CashV2ByString cmpCash; + private final CashV2ByIdAndString tcrCash; + + public ExecCurrCompCash() { + this.fixIncSecCash = add(new CashV2ByString<>("fixIncSecCash", Security::getSecuritySymbol)); + this.mmsCash = add(new CashV2ByString<>("mmsCash", Security::getSecuritySymbol)); + this.eqtySecCash = add(new CashV2ByString<>("eqtySecCash", Security::getSecuritySymbol)); + this.currPairSecCash = add(new CashV2ByString<>("currPairSecCash", Security::getSecuritySymbol)); + this.cmpCash = add(new CashV2ByString<>("cmpCash", Company::getTradingCode)); + this.tcrCash = add(new CashV2ByIdAndString<>("tcrCash", tcr -> new CashV2ByIdAndString.CustomKey(tcr.getCompanyId(), tcr.getCode()))); + } + } } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/execution/ExecutionDepositComponent.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/execution/ExecutionDepositComponent.java index 6cd27e1ad..45a94cf8a 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/execution/ExecutionDepositComponent.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/execution/ExecutionDepositComponent.java @@ -8,7 +8,6 @@ import java.util.Collection; import java.util.HashMap; import java.util.Map; import java.util.Optional; -import java.util.function.Function; import org.apache.kafka.clients.producer.Producer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -22,6 +21,7 @@ import ru.clearing.classes.statics.data.execution.ExecutionDeposit; import ru.clearing.classes.statics.data.misc.Listing; import ru.clearing.classes.statics.data.misc.STrades; import ru.clearing.classes.statics.data.registry.TradingClearingRegistry; +import ru.clearing.classes.statics.data.security.MoneyMarketSecurity; import ru.clearing.classes.statics.data.security.Security; import ru.spcex.clearing.error.ClearingError; import ru.spcex.clearing.error.ClearingException; @@ -31,6 +31,9 @@ import ru.spcex.clearing.platform.messaging.domain.cud.registry.DealRegisterNewR import ru.spcex.clearing.platform.messaging.domain.cud.registry.ExecutionType; import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.clearing.service.validation.ValidationStored; +import ru.spcex.clearing.service.validation.strades.CompanyByTradingCodeCashingValidationRule; +import ru.spcex.clearing.service.validation.strades.STradesSecurityPresentDepositValidationRule; +import ru.spcex.clearing.service.validation.strades.TcrByCodeAndCmpIsActiveCashingValidationRule; import ru.spcex.platform.enumeration.MoneyFlowSide; import ru.spcex.platform.enumeration.Section; import ru.spcex.platform.enumeration.Side; @@ -39,6 +42,11 @@ import ru.spcex.platform.imdg.api.ImdgId; 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.iml.hazelcast.adapter.CashCloser; +import ru.spcex.platform.imdg.iml.hazelcast.adapter.CashV2ByIdAndString; +import ru.spcex.platform.imdg.iml.hazelcast.adapter.CashV2ByString; +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.IEnumId; import ru.spcex.platform.utils.enumeration.IEnumKey; @@ -46,6 +54,7 @@ import ru.spcex.platform.utils.enumeration.IMessageResolver; import ru.spcex.platform.utils.log.ExceptionUtils; import ru.spcex.platform.utils.time.TimeUtil; import ru.spcex.platform.utils.validation.IValidator; +import ru.spcex.platform.utils.validation.ValidatorImpl; /** * 1.35. executionDeposit - Сделки @@ -59,29 +68,33 @@ public class ExecutionDepositComponent implements IExecutionUploadComponent { private final Imdg sTradeImdg; private final Imdg executionDepositImdg; private final Imdg listingImdg; + private final Imdg imdgMoneyMarketSecurity; + private final Imdg imdgCompany; + private final Imdg imdgTradingClearingRegistry; private final ImdgId idGen; //fixme ждать ТЗ Long tradeNum; Instant tradingDay; private final IMessageResolver msgResolver; - private final Function stradesValidator; private final KafkaSender kafkaSender; private final Map cash = new HashMap<>(); private static final DateTimeFormatter contractFormatter = DateTimeFormatter.ofPattern("ddMMyy"); + private final ExecDepoCompCash valCash = new ExecDepoCompCash(); @Autowired public ExecutionDepositComponent(ImdgProvider imdgProvider, Producer kafka, - @Qualifier("sTradesValidatorDeposit") Function stradesValidator, @Qualifier("kafkaSenderWithoutRequestInfo") KafkaSender kafkaSender, IMessageResolver msgResolver) { this.sTradeImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_STrades, STrades.class); this.executionDepositImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionDeposit, ExecutionDeposit.class); this.listingImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Listing, Listing.class); + this.imdgMoneyMarketSecurity = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_MoneyMarketSecurity, MoneyMarketSecurity.class, null); + this.imdgCompany = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_Company, Company.class, null); + this.imdgTradingClearingRegistry = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class, null); this.idGen = imdgProvider.getImdgIdGenerator(); - this.stradesValidator = stradesValidator; this.kafkaSender = kafkaSender; this.msgResolver = msgResolver; @@ -103,94 +116,96 @@ public class ExecutionDepositComponent implements IExecutionUploadComponent { } public void processNewTS(Long fromId) { - boolean cleanLoad = fromId == null; - if (cleanLoad) { - log.info("clearing cash..."); - cash.clear(); - initExecCash(); - } - log.debug("Start check new S_TRADE after {}, fromId={}", tradingDay, fromId); - //выбираем STrades на сегодня с правильным section - Collection sTrades; - if (cleanLoad) { - sTrades = sTradeImdg.getCollectionObjectsByFieldValues(Map.of( - "tradeDate", LocalDate.now(), - "section", Section.MKR.getKey() - )); - } else { - ImdgPredicateBuilder pb = sTradeImdg.predicateBuilder(); - ImdgPredicate filter = pb.and( - pb.equals("tradeDate", LocalDate.now()), - pb.equals("section", Section.MKR.getKey()), - pb.greatEqual("id", fromId) - ); - sTrades = sTradeImdg.getCollectionObjectsByPredicate(filter); - } - log.info("Found {} s_trade for today (fromId={})", sTrades.size(), fromId); - - if (sTrades.isEmpty()) { - logError(ClearingError.NewDealsNotFound); - return; - } - - //убираем уже добавленные в ExecutionDeposit - if (!cleanLoad) { - log.info("filtering strades based on ExecutionDeposit cash..."); - sTrades.removeIf(sTrd -> { - Side sTrdSide = IEnumKey.getEnumByKey(Side.class, sTrd.getOperation()); - if (sTrdSide == null || !(sTrdSide.equals(Side.BUY) || sTrdSide.equals(Side.SELL))) { - log.warn("cannot parse strade.tradeNum={} strade.operation {}", sTrd.getTradeNum(), sTrd.getOperation()); - return true; - } - MoneyFlowSide excDepSide = sTrdSide.equals(Side.BUY) ? MoneyFlowSide.BUY : MoneyFlowSide.SELL; - ExecUploadKey excKey = ExecUploadKey.cash(sTrd.getTradeDate(), - sTrd.getTradeNum(), - excDepSide.getKey()); - return cash.containsKey(excKey); - }); - } - log.info("{} strades left after already-added filtering", sTrades.size()); - - Map execsToInsert = new HashMap<>(); - for (STrades sTrd : sTrades) { - log.trace("S_TRADE[{}] new", sTrd.getId()); - - //проверки - IValidator validator = stradesValidator.apply(sTrd); - Optional error = validator.tillFirstError(); - if (error.isPresent()) { - logError(sTrd.getId(), error.get()); - continue; + try (valCash) { + boolean cleanLoad = fromId == null; + if (cleanLoad) { + log.info("clearing cash..."); + cash.clear(); + initExecCash(); } - ExecutionDeposit newED; - try { - newED = createExecutionDeposit(sTrd, validator); - ExecUploadKey execKey = ExecUploadKey.cash(newED); - ExecUploadCashInfo storedId; - if (!cleanLoad || (storedId = cash.get(execKey)) == null) { - newED.setId(idGen.nextId()); - cash.put(execKey, ExecUploadCashInfo.cash(newED)); - log.debug("New executionDeposit.id={} was created.", newED.getId()); - } else { - newED.setId(storedId.id()); - newED.setUpdated(newED.getCreated()); - newED.setCreated(storedId.createDt()); - log.debug("executionDeposit.id={} was updated.", newED.getId()); - } - execsToInsert.put(newED.getId(), newED); - //executionDepositImdg.insert(newED); - sendNotification(newED); - } catch (ClearingException ce) { - auditMessage(ce); - } catch (Exception e) { - log.error("When create new ExecutionDeposit by STrade[{}] error: {}", sTrd.getId(), ExceptionUtils.getStackTrace(e)); + log.debug("Start check new S_TRADE after {}, fromId={}", tradingDay, fromId); + //выбираем STrades на сегодня с правильным section + Collection sTrades; + if (cleanLoad) { + sTrades = sTradeImdg.getCollectionObjectsByFieldValues(Map.of( + "tradeDate", LocalDate.now(), + "section", Section.MKR.getKey() + )); + } else { + ImdgPredicateBuilder pb = sTradeImdg.predicateBuilder(); + ImdgPredicate filter = pb.and( + pb.equals("tradeDate", LocalDate.now()), + pb.equals("section", Section.MKR.getKey()), + pb.greatEqual("id", fromId) + ); + sTrades = sTradeImdg.getCollectionObjectsByPredicate(filter); } - } - log.info("batch putAll {} ExecutionDeposit", execsToInsert.size()); - executionDepositImdg.putAll(execsToInsert, 200); + log.info("Found {} s_trade for today (fromId={})", sTrades.size(), fromId); - Long newMaxTradeNum = sTrades.stream().mapToLong(STrades::getTradeNum).max().orElseGet(() -> tradeNum); - log.info("Process completed. Next tradeNum is {}", newMaxTradeNum); + if (sTrades.isEmpty()) { + logError(ClearingError.NewDealsNotFound); + return; + } + + //убираем уже добавленные в ExecutionDeposit + if (!cleanLoad) { + log.info("filtering strades based on ExecutionDeposit cash..."); + sTrades.removeIf(sTrd -> { + Side sTrdSide = IEnumKey.getEnumByKey(Side.class, sTrd.getOperation()); + if (sTrdSide == null || !(sTrdSide.equals(Side.BUY) || sTrdSide.equals(Side.SELL))) { + log.warn("cannot parse strade.tradeNum={} strade.operation {}", sTrd.getTradeNum(), sTrd.getOperation()); + return true; + } + MoneyFlowSide excDepSide = sTrdSide.equals(Side.BUY) ? MoneyFlowSide.BUY : MoneyFlowSide.SELL; + ExecUploadKey excKey = ExecUploadKey.cash(sTrd.getTradeDate(), + sTrd.getTradeNum(), + excDepSide.getKey()); + return cash.containsKey(excKey); + }); + } + log.info("{} strades left after already-added filtering", sTrades.size()); + + Map execsToInsert = new HashMap<>(); + for (STrades sTrd : sTrades) { + log.trace("S_TRADE[{}] new", sTrd.getId()); + + //проверки + IValidator validator = valFor(sTrd); + Optional error = validator.tillFirstError(); + if (error.isPresent()) { + logError(sTrd.getId(), error.get()); + continue; + } + ExecutionDeposit newED; + try { + newED = createExecutionDeposit(sTrd, validator); + ExecUploadKey execKey = ExecUploadKey.cash(newED); + ExecUploadCashInfo storedId; + if (!cleanLoad || (storedId = cash.get(execKey)) == null) { + newED.setId(idGen.nextId()); + cash.put(execKey, ExecUploadCashInfo.cash(newED)); + log.debug("New executionDeposit.id={} was created.", newED.getId()); + } else { + newED.setId(storedId.id()); + newED.setUpdated(newED.getCreated()); + newED.setCreated(storedId.createDt()); + log.debug("executionDeposit.id={} was updated.", newED.getId()); + } + execsToInsert.put(newED.getId(), newED); + //executionDepositImdg.insert(newED); + sendNotification(newED); + } catch (ClearingException ce) { + auditMessage(ce); + } catch (Exception e) { + log.error("When create new ExecutionDeposit by STrade[{}] error: {}", sTrd.getId(), ExceptionUtils.getStackTrace(e)); + } + } + log.info("batch putAll {} ExecutionDeposit", execsToInsert.size()); + executionDepositImdg.putAll(execsToInsert, 200); + + Long newMaxTradeNum = sTrades.stream().mapToLong(STrades::getTradeNum).max().orElseGet(() -> tradeNum); + log.info("Process completed. Next tradeNum is {}", newMaxTradeNum); + } } @@ -322,7 +337,7 @@ public class ExecutionDepositComponent implements IExecutionUploadComponent { } private Optional searchTcrByStrades(STrades sTrades) { - IValidator counterValidator = stradesValidator.apply(sTrades); + IValidator counterValidator = valFor(sTrades); Optional err = counterValidator.tillFirstError(); if (err.isPresent()) { log.error("couldn't extract setCounterPartyTradingClearingRegistryId from strades.id={}", sTrades.getId()); @@ -335,4 +350,37 @@ public class ExecutionDepositComponent implements IExecutionUploadComponent { public void initExecCash() { ExecutionUploadCashUtil.loadExecutions(executionDepositImdg, cash); } + + + public IValidator valFor(STrades sTrades) { + ImdgValidationContext context = new ImdgValidationContext<>(); + context.setValidatedObject(sTrades); + context.addImdg(IMDGDistributedNames.Map_MoneyMarketSecurity, imdgMoneyMarketSecurity); + context.addImdg(IMDGDistributedNames.Map_Company, imdgCompany); + context.addImdg(IMDGDistributedNames.Map_TradingClearingRegistry, imdgTradingClearingRegistry); + context.setLogPrefix(LogPrefixId.INSTANCE); + return new ValidatorImpl<>(context, + new STradesSecurityPresentDepositValidationRule(valCash.mmsCash), + new CompanyByTradingCodeCashingValidationRule<>( + STrades::getFirmId, valCash.cmpCash, ValidationStored.STradesCompany + ), + new CompanyByTradingCodeCashingValidationRule<>( + STrades::getCpFirmId, valCash.cmpCash, ValidationStored.STradesCompany + ), + new TcrByCodeAndCmpIsActiveCashingValidationRule(valCash.tcrCash) + ); + } + + private static class ExecDepoCompCash extends CashCloser { + private final CashV2ByString mmsCash; + private final CashV2ByString cmpCash; + private final CashV2ByIdAndString tcrCash; + + public ExecDepoCompCash() { + this.mmsCash = add(new CashV2ByString<>("mmsCash", Security::getSecuritySymbol)); + this.cmpCash = add(new CashV2ByString<>("cmpCash", Company::getTradingCode)); + this.tcrCash = add(new CashV2ByIdAndString<>("tcrCash", tcr -> new CashV2ByIdAndString.CustomKey(tcr.getCompanyId(), tcr.getCode()))); + } + } + } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/execution/ExecutionFondComponent.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/execution/ExecutionFondComponent.java index 58a041b59..bc5cdafec 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/execution/ExecutionFondComponent.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/execution/ExecutionFondComponent.java @@ -8,8 +8,6 @@ import java.util.Collection; import java.util.HashMap; import java.util.Map; import java.util.Optional; -import java.util.function.Function; -import org.apache.kafka.clients.producer.Producer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Autowired; @@ -20,6 +18,7 @@ import org.springframework.stereotype.Component; import ru.clearing.classes.statics.data.account.ClientCode; import ru.clearing.classes.statics.data.company.Company; import ru.clearing.classes.statics.data.execution.ExecutionFond; +import ru.clearing.classes.statics.data.instrument.issue.EquitySecurity; import ru.clearing.classes.statics.data.instrument.issue.FixedIncomeCashFlow; import ru.clearing.classes.statics.data.instrument.issue.FixedIncomeSecurity; import ru.clearing.classes.statics.data.misc.Listing; @@ -27,6 +26,8 @@ import ru.clearing.classes.statics.data.misc.Market; import ru.clearing.classes.statics.data.misc.SCrossRate; import ru.clearing.classes.statics.data.misc.STrades; import ru.clearing.classes.statics.data.registry.TradingClearingRegistry; +import ru.clearing.classes.statics.data.security.CurrencyPairSecurity; +import ru.clearing.classes.statics.data.security.MoneyMarketSecurity; import ru.clearing.classes.statics.data.security.Security; import ru.spcex.clearing.config.element.ClearingServiceSettings; import ru.spcex.clearing.error.ClearingError; @@ -38,6 +39,9 @@ import ru.spcex.clearing.platform.messaging.domain.cud.registry.DealRegisterNewR import ru.spcex.clearing.platform.messaging.domain.cud.registry.ExecutionType; import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.clearing.service.validation.ValidationStored; +import ru.spcex.clearing.service.validation.strades.CompanyByTradingCodeCashingValidationRule; +import ru.spcex.clearing.service.validation.strades.STradesSecurityPresentValidationRule; +import ru.spcex.clearing.service.validation.strades.TcrByCodeAndCmpIsActiveCashingValidationRule; import ru.spcex.platform.enumeration.CurrencyCode; import ru.spcex.platform.enumeration.InstrumentType; import ru.spcex.platform.enumeration.ObjectType; @@ -49,6 +53,11 @@ import ru.spcex.platform.imdg.api.ImdgId; 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.iml.hazelcast.adapter.CashCloser; +import ru.spcex.platform.imdg.iml.hazelcast.adapter.CashV2ByIdAndString; +import ru.spcex.platform.imdg.iml.hazelcast.adapter.CashV2ByString; +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.IEnumId; import ru.spcex.platform.utils.enumeration.IEnumKey; @@ -59,6 +68,7 @@ import static ru.spcex.platform.utils.number.BigDecimalUtil.safeBD; import ru.spcex.platform.utils.text.TextUtil; import ru.spcex.platform.utils.time.TimeUtil; import ru.spcex.platform.utils.validation.IValidator; +import ru.spcex.platform.utils.validation.ValidatorImpl; @Component @EnableScheduling @@ -72,24 +82,28 @@ public class ExecutionFondComponent implements IExecutionUploadComponent { private final Imdg fixedIncomeCashFlowImdg; private final Imdg marketImdg; private final Imdg crossRateImdg; + private final Imdg currencyPairSecurityImdg; + private final Imdg imdgMoneyMarketSecurity; + private final Imdg fixedIncomeSecurityImdg; + private final Imdg equitySecurityImdg; + private final Imdg imdgCompany; + private final Imdg imdgTradingClearingRegistry; private final ImdgId idGen; //fixme ждать ТЗ Long tradeNum; Instant tradingDay; private final IMessageResolver msgResolver = new SimpleMessageResolver(); - private final Function stradesValidator; private final KafkaSender kafkaSender; private final NotificationSender notifications; private final boolean valuation; private final Map cash = new HashMap<>(); + private final ExecFondCompCash valCash = new ExecFondCompCash(); @Autowired public ExecutionFondComponent(ImdgProvider imdgProvider, ClearingServiceSettings settings, - Producer kafka, - @Qualifier("sTradesValidatorFond") Function stradesValidator, @Qualifier("kafkaSenderWithoutRequestInfo") KafkaSender kafkaSender, NotificationSender notifications) { this.sTradeImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_STrades, STrades.class); this.executionFondImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionFond, ExecutionFond.class); @@ -98,9 +112,15 @@ public class ExecutionFondComponent implements IExecutionUploadComponent { this.marketImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Market, Market.class); this.fixedIncomeCashFlowImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_FixedIncomeCashFlow, FixedIncomeCashFlow.class); this.crossRateImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SCrossRate, SCrossRate.class); + + this.currencyPairSecurityImdg = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_CurrencyPairSecurity, CurrencyPairSecurity.class, null); + this.imdgMoneyMarketSecurity = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_MoneyMarketSecurity, MoneyMarketSecurity.class, null); + this.fixedIncomeSecurityImdg = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_FixedIncomeSecurity, FixedIncomeSecurity.class, null); + this.equitySecurityImdg = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_EquitySecurity, EquitySecurity.class, null); + this.imdgCompany = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_Company, Company.class, null); + this.imdgTradingClearingRegistry = imdgProvider.getCashingImdg(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class, null); this.idGen = imdgProvider.getImdgIdGenerator(); this.valuation = settings.getTrade().getValuation(); //fixme - this.stradesValidator = stradesValidator; this.kafkaSender = kafkaSender; this.notifications = notifications; @@ -122,94 +142,96 @@ public class ExecutionFondComponent implements IExecutionUploadComponent { } public void processNewTS(Long fromId) { - boolean cleanLoad = fromId == null; - if (cleanLoad) { - log.info("clearing cash..."); - cash.clear(); - initExecCash(); - } - log.debug("Start check new S_TRADE after {}, fromId={}", tradingDay, fromId); - //выбираем STrades на сегодня с правильным section - Collection sTrades; - if (cleanLoad) { - sTrades = sTradeImdg.getCollectionObjectsByFieldValues(Map.of( - "tradeDate", LocalDate.now(), - "section", Section.FOND.getKey() - )); - } else { - ImdgPredicateBuilder pb = sTradeImdg.predicateBuilder(); - ImdgPredicate filter = pb.and( - pb.equals("tradeDate", LocalDate.now()), - pb.equals("section", Section.FOND.getKey()), - pb.greatEqual("id", fromId) - ); - sTrades = sTradeImdg.getCollectionObjectsByPredicate(filter); - } - log.info("Found {} s_trade for today (fromId={})", sTrades.size(), fromId); - - if (sTrades.isEmpty()) { - logError(ClearingError.NewDealsNotFound); - return; - } - - //убираем уже добавленные в ExecutionDeposit - if (!cleanLoad) { - log.info("filtering strades based on ExecutionFond cash..."); - sTrades.removeIf(sTrd -> { - Side sTrdSide = IEnumKey.getEnumByKey(Side.class, sTrd.getOperation()); - if (sTrdSide == null || !(sTrdSide.equals(Side.BUY) || sTrdSide.equals(Side.SELL))) { - log.warn("cannot parse strade.tradeNum={} strade.operation {}", sTrd.getTradeNum(), sTrd.getOperation()); - return true; - } - Side excDepSide = sTrdSide.equals(Side.BUY) ? Side.BUY : Side.SELL; - ExecUploadKey excKey = ExecUploadKey.cash(sTrd.getTradeDate(), - sTrd.getTradeNum(), - excDepSide.getKey()); - return cash.containsKey(excKey); - }); - } - log.info("{} strades left after already-added filtering", sTrades.size()); - - Map execsToInsert = new HashMap<>(); - for (STrades sTrd : sTrades) { - log.info("new S_TRADE[{}], valuation {}", sTrd.getId(), valuation); - - //проверки - IValidator validator = stradesValidator.apply(sTrd); - Optional error = validator.tillFirstError(); - if (error.isPresent()) { - logError(sTrd.getId(), error.get()); - continue; + try (valCash) { + boolean cleanLoad = fromId == null; + if (cleanLoad) { + log.info("clearing cash..."); + cash.clear(); + initExecCash(); } - ExecutionFond newED; - try { - newED = createExecutionFond(sTrd, validator); - ExecUploadKey execKey = ExecUploadKey.cash(newED); - ExecUploadCashInfo storedId; - if (!cleanLoad || (storedId = cash.get(execKey)) == null) { - newED.setId(idGen.nextId()); - cash.put(execKey, ExecUploadCashInfo.cash(newED)); - log.debug("New executionFond.id={} was created.", newED.getId()); - } else { - newED.setId(storedId.id()); - newED.setUpdated(newED.getCreated()); - newED.setCreated(storedId.createDt()); - log.debug("executionFond.id={} was updated.", newED.getId()); - } - execsToInsert.put(newED.getId(), newED); - //executionFondImdg.insert(newED); - sendNotification(newED); - } catch (ClearingException ce) { - auditMessage(ce); - } catch (Exception e) { - log.error("When create new ExecutionFond by STrade[{}] error: {}", sTrd.getId(), ExceptionUtils.getStackTrace(e)); + log.debug("Start check new S_TRADE after {}, fromId={}", tradingDay, fromId); + //выбираем STrades на сегодня с правильным section + Collection sTrades; + if (cleanLoad) { + sTrades = sTradeImdg.getCollectionObjectsByFieldValues(Map.of( + "tradeDate", LocalDate.now(), + "section", Section.FOND.getKey() + )); + } else { + ImdgPredicateBuilder pb = sTradeImdg.predicateBuilder(); + ImdgPredicate filter = pb.and( + pb.equals("tradeDate", LocalDate.now()), + pb.equals("section", Section.FOND.getKey()), + pb.greatEqual("id", fromId) + ); + sTrades = sTradeImdg.getCollectionObjectsByPredicate(filter); } - } - log.info("batch putAll {} ExecutionFond", execsToInsert.size()); - executionFondImdg.putAll(execsToInsert, 200); + log.info("Found {} s_trade for today (fromId={})", sTrades.size(), fromId); - Long newMaxTradeNum = sTrades.stream().mapToLong(STrades::getTradeNum).max().orElseGet(() -> tradeNum); - log.info("Process completed. Next tradeNum is {}", newMaxTradeNum); + if (sTrades.isEmpty()) { + logError(ClearingError.NewDealsNotFound); + return; + } + + //убираем уже добавленные в ExecutionDeposit + if (!cleanLoad) { + log.info("filtering strades based on ExecutionFond cash..."); + sTrades.removeIf(sTrd -> { + Side sTrdSide = IEnumKey.getEnumByKey(Side.class, sTrd.getOperation()); + if (sTrdSide == null || !(sTrdSide.equals(Side.BUY) || sTrdSide.equals(Side.SELL))) { + log.warn("cannot parse strade.tradeNum={} strade.operation {}", sTrd.getTradeNum(), sTrd.getOperation()); + return true; + } + Side excDepSide = sTrdSide.equals(Side.BUY) ? Side.BUY : Side.SELL; + ExecUploadKey excKey = ExecUploadKey.cash(sTrd.getTradeDate(), + sTrd.getTradeNum(), + excDepSide.getKey()); + return cash.containsKey(excKey); + }); + } + log.info("{} strades left after already-added filtering", sTrades.size()); + + Map execsToInsert = new HashMap<>(); + for (STrades sTrd : sTrades) { + log.info("new S_TRADE[{}], valuation {}", sTrd.getId(), valuation); + + //проверки + IValidator validator = valFor(sTrd); + Optional error = validator.tillFirstError(); + if (error.isPresent()) { + logError(sTrd.getId(), error.get()); + continue; + } + ExecutionFond newED; + try { + newED = createExecutionFond(sTrd, validator); + ExecUploadKey execKey = ExecUploadKey.cash(newED); + ExecUploadCashInfo storedId; + if (!cleanLoad || (storedId = cash.get(execKey)) == null) { + newED.setId(idGen.nextId()); + cash.put(execKey, ExecUploadCashInfo.cash(newED)); + log.debug("New executionFond.id={} was created.", newED.getId()); + } else { + newED.setId(storedId.id()); + newED.setUpdated(newED.getCreated()); + newED.setCreated(storedId.createDt()); + log.debug("executionFond.id={} was updated.", newED.getId()); + } + execsToInsert.put(newED.getId(), newED); + //executionFondImdg.insert(newED); + sendNotification(newED); + } catch (ClearingException ce) { + auditMessage(ce); + } catch (Exception e) { + log.error("When create new ExecutionFond by STrade[{}] error: {}", sTrd.getId(), ExceptionUtils.getStackTrace(e)); + } + } + log.info("batch putAll {} ExecutionFond", execsToInsert.size()); + executionFondImdg.putAll(execsToInsert, 200); + + Long newMaxTradeNum = sTrades.stream().mapToLong(STrades::getTradeNum).max().orElseGet(() -> tradeNum); + log.info("Process completed. Next tradeNum is {}", newMaxTradeNum); + } } @@ -374,4 +396,48 @@ public class ExecutionFondComponent implements IExecutionUploadComponent { public void initExecCash() { ExecutionUploadCashUtil.loadExecutions(executionFondImdg, cash); } + + + public IValidator valFor(STrades sTrades) { + ImdgValidationContext context = new ImdgValidationContext<>(); + context.setValidatedObject(sTrades); + context.addImdg(IMDGDistributedNames.Map_CurrencyPairSecurity, currencyPairSecurityImdg); + context.addImdg(IMDGDistributedNames.Map_MoneyMarketSecurity, imdgMoneyMarketSecurity); + context.addImdg(IMDGDistributedNames.Map_FixedIncomeSecurity, fixedIncomeSecurityImdg); + context.addImdg(IMDGDistributedNames.Map_EquitySecurity, equitySecurityImdg); + context.addImdg(IMDGDistributedNames.Map_Company, imdgCompany); + context.addImdg(IMDGDistributedNames.Map_TradingClearingRegistry, imdgTradingClearingRegistry); + context.setLogPrefix(LogPrefixId.INSTANCE); + return new ValidatorImpl<>(context, + new STradesSecurityPresentValidationRule(valCash.fixIncSecCash, + valCash.mmsCash, + valCash.eqtySecCash, + valCash.currPairSecCash), + new CompanyByTradingCodeCashingValidationRule<>( + STrades::getFirmId, valCash.cmpCash, ValidationStored.STradesCompany + ), + new CompanyByTradingCodeCashingValidationRule<>( + STrades::getCpFirmId, valCash.cmpCash, ValidationStored.STradesCompany + ), + new TcrByCodeAndCmpIsActiveCashingValidationRule(valCash.tcrCash) + ); + } + + private static class ExecFondCompCash extends CashCloser { + private final CashV2ByString fixIncSecCash; + private final CashV2ByString mmsCash; + private final CashV2ByString eqtySecCash; + private final CashV2ByString currPairSecCash; + private final CashV2ByString cmpCash; + private final CashV2ByIdAndString tcrCash; + + public ExecFondCompCash() { + this.fixIncSecCash = add(new CashV2ByString<>("fixIncSecCash", Security::getSecuritySymbol)); + this.mmsCash = add(new CashV2ByString<>("mmsCash", Security::getSecuritySymbol)); + this.eqtySecCash = add(new CashV2ByString<>("eqtySecCash", Security::getSecuritySymbol)); + this.currPairSecCash = add(new CashV2ByString<>("currPairSecCash", Security::getSecuritySymbol)); + this.cmpCash = add(new CashV2ByString<>("cmpCash", Company::getTradingCode)); + this.tcrCash = add(new CashV2ByIdAndString<>("tcrCash", tcr -> new CashV2ByIdAndString.CustomKey(tcr.getCompanyId(), tcr.getCode()))); + } + } } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/strades/CompanyByTradingCodeCashingValidationRule.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/strades/CompanyByTradingCodeCashingValidationRule.java new file mode 100644 index 000000000..492add30f --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/strades/CompanyByTradingCodeCashingValidationRule.java @@ -0,0 +1,66 @@ +package ru.spcex.clearing.service.validation.strades; + +import java.util.Optional; +import java.util.function.Function; +import ru.clearing.classes.statics.data.company.Company; +import ru.spcex.clearing.error.ClearingError; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.predicate.ImdgPredicate; +import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder; +import ru.spcex.platform.imdg.iml.hazelcast.adapter.CashV2ByString; +import ru.spcex.platform.imdg.validation.ImdgValidationContext; +import ru.spcex.platform.utils.enumeration.EnumMessage; +import ru.spcex.platform.utils.text.TextUtil; +import ru.spcex.platform.utils.validation.IValidationRule; + +public class CompanyByTradingCodeCashingValidationRule + implements IValidationRule> { + + private final CashV2ByString cmpCash; + private final Function extractor; + private final Enum storedKey; + + + public CompanyByTradingCodeCashingValidationRule(Function extractor, + CashV2ByString cmpCash, + Enum storedKey) { + this.cmpCash = cmpCash; + this.extractor = extractor; + this.storedKey = storedKey; + } + + public CompanyByTradingCodeCashingValidationRule(Function extractor, + CashV2ByString cmpCash) { + this.cmpCash = cmpCash; + this.extractor = extractor; + this.storedKey = null; + } + + @Override + public String ruleName() { + return "CompanyByTradingCodeCashingValidationRule"; + } + + @Override + public Optional validate(ImdgValidationContext context) { + T validatedObject = context.getValidatedObject(); + String tradingCode = extractor.apply(validatedObject); + if (TextUtil.isEmpty(tradingCode)) { + return of(ClearingError.CompanyNotFound, tradingCode); + } + Imdg companyImdg = context.obtainMap(IMDGDistributedNames.Map_Company, Company.class); + ImdgPredicateBuilder pb = companyImdg.predicateBuilder(); + ImdgPredicate prdct = pb.cashed( + pb.equals("tradingCode", tradingCode), cmpCash, tradingCode + ); + Company company = companyImdg.getFirstObjectByPredicate(prdct); + if (company == null) { + return of(ClearingError.CompanyNotFound, tradingCode); + } + if (storedKey != null) { + context.storeObject(storedKey, company); + } + return empty(); + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/strades/STradesSecurityPresentDepositValidationRule.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/strades/STradesSecurityPresentDepositValidationRule.java new file mode 100644 index 000000000..d1f4aa8e5 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/strades/STradesSecurityPresentDepositValidationRule.java @@ -0,0 +1,52 @@ +package ru.spcex.clearing.service.validation.strades; + +import java.util.Optional; +import ru.clearing.classes.statics.data.misc.STrades; +import ru.clearing.classes.statics.data.security.MoneyMarketSecurity; +import ru.clearing.classes.statics.data.security.Security; +import ru.spcex.clearing.error.ClearingError; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.service.validation.ValidationStored; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.predicate.ImdgPredicate; +import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder; +import ru.spcex.platform.imdg.iml.hazelcast.adapter.CashV2ByString; +import ru.spcex.platform.imdg.validation.ImdgValidationContext; +import ru.spcex.platform.utils.enumeration.EnumMessage; +import ru.spcex.platform.utils.text.TextUtil; +import ru.spcex.platform.utils.validation.IValidationRule; + +public class STradesSecurityPresentDepositValidationRule implements IValidationRule> { + + private final CashV2ByString mmsCash; + + public STradesSecurityPresentDepositValidationRule(CashV2ByString mmsCash) { + this.mmsCash = mmsCash; + } + + @Override + public String ruleName() { + return "STradesSecurityPresentCurrencyPairValidationRule"; + } + + @Override + public Optional validate(ImdgValidationContext context) { + + STrades validatedObject = context.getValidatedObject(); + if (TextUtil.isEmpty(validatedObject.getSecCode())) { + return of(ClearingError.SecurityNotFound, validatedObject.getSecCode()); + } + String secCode = validatedObject.getSecCode().trim(); + secCode = secCode.substring(0, Math.min(7, secCode.length())); + Imdg securityImdg = context.obtainMap(IMDGDistributedNames.Map_MoneyMarketSecurity, MoneyMarketSecurity.class); + ImdgPredicateBuilder pb = securityImdg.predicateBuilder(); + ImdgPredicate prdct = pb.equals("securitySymbol", secCode); + prdct = pb.cashed(prdct, mmsCash, secCode); + Security security = securityImdg.getFirstObjectByPredicate(prdct); + if (security == null) { + return of(ClearingError.SecurityNotFound, secCode); + } + context.storeObject(ValidationStored.STradesSecurity, security); + return empty(); + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/strades/STradesSecurityPresentValidationRule.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/strades/STradesSecurityPresentValidationRule.java new file mode 100644 index 000000000..da808f255 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/strades/STradesSecurityPresentValidationRule.java @@ -0,0 +1,57 @@ +package ru.spcex.clearing.service.validation.strades; + +import java.util.Optional; +import ru.clearing.classes.statics.data.misc.STrades; +import ru.clearing.classes.statics.data.security.Security; +import ru.spcex.clearing.error.ClearingError; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.service.validation.ValidationStored; +import ru.spcex.platform.imdg.api.predicate.specific.SecuritySelector; +import ru.spcex.platform.imdg.iml.hazelcast.adapter.CashV2ByString; +import ru.spcex.platform.imdg.validation.ImdgValidationContext; +import ru.spcex.platform.utils.enumeration.EnumMessage; +import ru.spcex.platform.utils.text.TextUtil; +import ru.spcex.platform.utils.validation.IValidationRule; + +public class STradesSecurityPresentValidationRule implements IValidationRule> { + + private final CashV2ByString fixedIncomeSecurityCash; + private final CashV2ByString moneyMarketSecurityCash; + private final CashV2ByString equitySecurityCash; + private final CashV2ByString currencyPairSecurityCash; + + public STradesSecurityPresentValidationRule(CashV2ByString fixedIncomeSecurityCash, CashV2ByString moneyMarketSecurityCash, CashV2ByString equitySecurityCash, CashV2ByString currencyPairSecurityCash) { + this.fixedIncomeSecurityCash = fixedIncomeSecurityCash; + this.moneyMarketSecurityCash = moneyMarketSecurityCash; + this.equitySecurityCash = equitySecurityCash; + this.currencyPairSecurityCash = currencyPairSecurityCash; + } + + @Override + public String ruleName() { + return "STradesSecurityPresentCurrencyPairValidationRule"; + } + + @Override + public Optional validate(ImdgValidationContext context) { + STrades validatedObject = context.getValidatedObject(); + if (TextUtil.isEmpty(validatedObject.getSecCode())) { + return of(ClearingError.SecurityNotFound, validatedObject.getSecCode()); + } + + SecuritySelector slctr = new SecuritySelector<>( + context.obtainMap(IMDGDistributedNames.Map_FixedIncomeSecurity, Security.class), + context.obtainMap(IMDGDistributedNames.Map_MoneyMarketSecurity, Security.class), + context.obtainMap(IMDGDistributedNames.Map_EquitySecurity, Security.class), + context.obtainMap(IMDGDistributedNames.Map_CurrencyPairSecurity, Security.class) + ); + slctr.setCash(fixedIncomeSecurityCash, moneyMarketSecurityCash, equitySecurityCash, currencyPairSecurityCash); + String secCode = validatedObject.getSecCode().trim(); + Security security = slctr.selectSecurityBySymbolCashed(secCode); + if (security == null) { + return of(ClearingError.SecurityNotFound, secCode); + } + context.storeObject(ValidationStored.STradesSecurity, security); + return empty(); + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/strades/TcrByCodeAndCmpIsActiveCashingValidationRule.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/strades/TcrByCodeAndCmpIsActiveCashingValidationRule.java new file mode 100644 index 000000000..6c1048ddb --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/strades/TcrByCodeAndCmpIsActiveCashingValidationRule.java @@ -0,0 +1,71 @@ +package ru.spcex.clearing.service.validation.strades; + +import java.util.Optional; +import ru.clearing.classes.statics.data.company.Company; +import ru.clearing.classes.statics.data.misc.STrades; +import ru.clearing.classes.statics.data.registry.TradingClearingRegistry; +import ru.spcex.clearing.error.ClearingError; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.service.validation.ValidationStored; +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.api.predicate.ImdgPredicateBuilder; +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.text.TextUtil; +import ru.spcex.platform.utils.validation.IValidationRule; + +public class TcrByCodeAndCmpIsActiveCashingValidationRule implements IValidationRule> { + private final CashV2ByIdAndString tcrCash; + + public TcrByCodeAndCmpIsActiveCashingValidationRule(CashV2ByIdAndString tcrCash) { + this.tcrCash = tcrCash; + } + + @Override + public Optional validate(ImdgValidationContext context) { + STrades validatedObject = context.getValidatedObject(); + if (TextUtil.isEmpty(validatedObject.getAccount())) { + return of(ClearingError.TradingClearingRegistryNotFound, "", validatedObject.getAccount()); + } + if (context.getStoredObject(ValidationStored.STradesCompany) == null) { //maybe unnecessary + return of(ClearingError.CompanyNotFound, validatedObject.getSecCode()); + } + Imdg tcrImdg = context.obtainMap( + IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class + ); + + ImdgPredicateBuilder pb = tcrImdg.predicateBuilder(); + String code = validatedObject.getAccount().trim(); + Long cmpId = ((Company) context.getStoredObject(ValidationStored.STradesCompany)).getId(); + ImdgPredicate prdct = pb.and( + pb.equals("code", code), + pb.equals("companyId", cmpId) + ); + prdct = pb.cashed(prdct, tcrCash, new CashV2ByIdAndString.CustomKey(cmpId, code)); + TradingClearingRegistry tcr = tcrImdg.getFirstObjectByPredicate( + prdct + ); + + if (tcr == null) { + return of(ClearingError.TradingClearingRegistryNotFound, "", validatedObject.getAccount()); + } + // 2.4. Проверить, что найденный торгово-клиринговый регистр в активном состоянии: + // tradingClearingRegistry.status≠BLKD/SSPD/CLOS (см. справочник serviceStatus). + // Иначе записать в лог ошибку (5419) "Торгово-клиринговый регистр %s неактивен". + ServiceStatus status = IEnumKey.getEnumByKey(ServiceStatus.class, tcr.getStatus()); + if (IEnumKey.contains(status, ServiceStatus.Blocked, ServiceStatus.Suspended, ServiceStatus.Closed)) { + return of(ClearingError.TradingClearingRegistryNotActive, validatedObject.getAccount()); + } + context.storeObject(ValidationStored.STradesTradingClearingRegistry, tcr); + return empty(); + } + + @Override + public String ruleName() { + return "TcrByCodeAndCmpIsActiveCashingValidationRule"; + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/FormingPaymentInstructionAssets.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/FormingPaymentInstructionAssets.java index 0b17bf4b7..6cfbcfab2 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/FormingPaymentInstructionAssets.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/FormingPaymentInstructionAssets.java @@ -182,6 +182,7 @@ public class FormingPaymentInstructionAssets implements ISessionStage { for (Registry registry : AMBregistries) { Account tranAcc = accountImdg.getFirstObjectBySQL(("accountType = '%s' " + "and status = '%s' " + + //fixme empty for RUB "and currency = '%s' " + "and processingSign = '%s'") .formatted(AccountType.Tran.getKey(), diff --git a/platform-parent/platform-imdg-api/src/main/java/ru/spcex/platform/imdg/api/predicate/specific/SecuritySelector.java b/platform-parent/platform-imdg-api/src/main/java/ru/spcex/platform/imdg/api/predicate/specific/SecuritySelector.java index 31d4f5295..86bfa9ea7 100644 --- a/platform-parent/platform-imdg-api/src/main/java/ru/spcex/platform/imdg/api/predicate/specific/SecuritySelector.java +++ b/platform-parent/platform-imdg-api/src/main/java/ru/spcex/platform/imdg/api/predicate/specific/SecuritySelector.java @@ -1,19 +1,26 @@ package ru.spcex.platform.imdg.api.predicate.specific; +import java.util.Map; +import java.util.function.BiFunction; +import java.util.function.Function; import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.platform.classes.base.SpcexObjectBase; import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.ImdgProvider; - -import java.util.Map; -import java.util.function.Function; +import ru.spcex.platform.imdg.api.predicate.ImdgPredicate; +import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder; +import ru.spcex.platform.imdg.iml.hazelcast.adapter.CashV2ByString; public class SecuritySelector { private final Imdg fixedIncomeSecurityImdg; + private CashV2ByString fixedIncomeSecurityCash = null; private final Imdg moneyMarketSecurityImdg; + private CashV2ByString moneyMarketSecurityCash = null; private final Imdg equitySecurity; + private CashV2ByString equitySecurityCash = null; private final Imdg currencyPairSecurity; + private CashV2ByString currencyPairSecurityCash = null; public SecuritySelector(ImdgProvider imdgProvider, Class clazz) { this.fixedIncomeSecurityImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_FixedIncomeSecurity, clazz); @@ -32,6 +39,16 @@ public class SecuritySelector { this.currencyPairSecurity = currencyPairSecurity; } + public void setCash(CashV2ByString fixedIncomeSecurityCash, + CashV2ByString moneyMarketSecurityCash, + CashV2ByString equitySecuCash, + CashV2ByString currencyPairSecuCash) { + this.fixedIncomeSecurityCash = fixedIncomeSecurityCash; + this.moneyMarketSecurityCash = moneyMarketSecurityCash; + this.equitySecurityCash = equitySecuCash; + this.currencyPairSecurityCash = currencyPairSecuCash; + } + public T selectSecurityById(Long securityId) { return selectSecurityByAnything(imdg -> imdg.getSingleObjectByID(securityId)); } @@ -41,6 +58,14 @@ public class SecuritySelector { return selectSecurityByAnything(imdg -> imdg.getFirstObjectByFieldValues(fields)); } + public T selectSecurityBySymbolCashed(String securitySymbol) { + return selectSecurityByAnything((imdg, cash) -> { + ImdgPredicateBuilder pb = imdg.predicateBuilder(); + ImdgPredicate prdct = pb.equals("securitySymbol", securitySymbol); + return imdg.getFirstObjectByPredicate(pb.cashed(prdct, cash, securitySymbol)); + }); + } + private T selectSecurityByAnything(Function, T> retriever) { T security = retriever.apply(fixedIncomeSecurityImdg); if (security != null) return security; @@ -51,4 +76,15 @@ public class SecuritySelector { security = retriever.apply(currencyPairSecurity); return security;// null } + + private T selectSecurityByAnything(BiFunction, CashV2ByString, T> retriever) { + T security = retriever.apply(fixedIncomeSecurityImdg, fixedIncomeSecurityCash); + if (security != null) return security; + security = retriever.apply(moneyMarketSecurityImdg, moneyMarketSecurityCash); + if (security != null) return security; + security = retriever.apply(equitySecurity, equitySecurityCash); + if (security != null) return security; + security = retriever.apply(currencyPairSecurity, currencyPairSecurityCash); + return security; + } }