diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/error/ClearingError.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/error/ClearingError.java index d0ce0b9c6..88825fa8a 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/error/ClearingError.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/error/ClearingError.java @@ -50,6 +50,7 @@ public enum ClearingError implements IErrorEnumId { WrongField(5004L), GatewayTimeout(6003L), GatewayNotApproved(6004L), + CounterStradeNotFound(-1L) ; private final Long id; diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/execution/ExecutionAbstractComponent.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/execution/ExecutionAbstractComponent.java new file mode 100644 index 000000000..83110ab5e --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/execution/ExecutionAbstractComponent.java @@ -0,0 +1,69 @@ +package ru.spcex.clearing.service.execution; + +import java.util.Collection; +import org.slf4j.Logger; +import org.slf4j.event.Level; +import ru.spcex.clearing.notification.NotificationSender; +import ru.spcex.platform.enumeration.ObjectType; +import ru.spcex.platform.enumeration.Priority; +import ru.spcex.platform.utils.enumeration.EnumMessage; +import ru.spcex.platform.utils.enumeration.IEnumId; +import ru.spcex.platform.utils.enumeration.IMessageResolver; + +public abstract class ExecutionAbstractComponent implements IExecutionUploadComponent { + protected static final int EXTRA_ITERATIONS_WARN = 500; + protected static final int EXTRA_ITERATIONS_ERROR = 1000; + protected final NotificationSender notifications; + private final IMessageResolver msgResolver; + + + protected ExecutionAbstractComponent(NotificationSender notifications, IMessageResolver msgResolver) { + this.notifications = notifications; + this.msgResolver = msgResolver; + } + + protected abstract Logger log(); + + protected int iterationsWarn(Collection items) { + return items.size() + EXTRA_ITERATIONS_WARN; + } + + protected int iterationsError(Collection items) { + return items.size() + EXTRA_ITERATIONS_ERROR; + } + + protected boolean cycleExceptionalInterrupt(int i, int iterationWarn, int iterationErr) { + if (i == iterationWarn) { + logNotifyE("Высокое число сделок без пары"); + } + if (i > iterationErr) { + logNotifyE("Превышено число сделок без пары"); + log().error("FATAL: too many unpaired STrades. breaking the cycle"); + return true; + } + return false; + } + + protected void logError(IEnumId subject, Object... args) { + String err = msgResolver.resolve(new EnumMessage(subject, args)); + log().info("{}", err); +// notifications.sendNotification(ObjectType.vfrs, err, Priority.HIGH); + } + + protected void logNotifyW(String msg) { + logNotify(msg, Level.WARN); + } + + protected void logNotifyE(String msg) { + logNotify(msg, Level.ERROR); + } + + protected void logNotify(String msg, Level lvl) { + switch (lvl) { + case WARN -> log().warn(msg); + case ERROR -> log().error(msg); + default -> log().error("lvl not supported {}. {}", lvl, msg); + } + notifications.sendNotification(ObjectType.vfrs, msg, Priority.HIGH); + } +} 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 10f6cd517..78583b9b7 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 @@ -6,15 +6,14 @@ import java.time.Instant; import java.time.LocalDate; import java.time.format.DateTimeFormatter; 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 static java.lang.String.format; import org.apache.kafka.clients.producer.Producer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import org.slf4j.event.Level; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.scheduling.annotation.EnableScheduling; @@ -42,8 +41,9 @@ 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.CounterStradePresentValidationRule; import ru.spcex.clearing.service.validation.strades.STradesSecurityPresentValidationRule; -import ru.spcex.clearing.service.validation.strades.TcrByCodeAndCmpIsActiveCashingValidationRule; +import ru.spcex.clearing.service.validation.strades.TcrByCodeAndCmpIsActiveCashingValidationRuleV2; import ru.spcex.clearing.validation.common.rules.FieldNotBlankRequiredRule; import ru.spcex.platform.enumeration.MarketType; import ru.spcex.platform.enumeration.ObjectType; @@ -60,8 +60,8 @@ 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.collection.CollectionUtil; import ru.spcex.platform.utils.enumeration.EnumMessage; -import ru.spcex.platform.utils.enumeration.IEnumId; import ru.spcex.platform.utils.enumeration.IEnumKey; import ru.spcex.platform.utils.enumeration.IMessageResolver; import ru.spcex.platform.utils.log.ExceptionUtils; @@ -76,7 +76,7 @@ import ru.spcex.platform.utils.validation.ValidatorImpl; */ @Component @EnableScheduling -public class ExecutionCurrencyComponent implements IExecutionUploadComponent { +public class ExecutionCurrencyComponent extends ExecutionAbstractComponent { private final Logger log = LoggerFactory.getLogger(getClass()); private final Imdg sTradeImdg; @@ -109,6 +109,7 @@ public class ExecutionCurrencyComponent implements IExecutionUploadComponent { IMessageResolver msgResolver, NotificationSender notifications, ClearingServiceSettings serviceSettings) { + super(notifications, 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); @@ -133,7 +134,12 @@ public class ExecutionCurrencyComponent implements IExecutionUploadComponent { ExecutionUploadCashUtil.loadExecutions(executionCurrencyImdg, cash); } -// /** + @Override + protected Logger log() { + return log; + } + + // /** // * Сбрасывать каждый день в 01:00:01 "0 1 0 1 * ?" // */ // @Scheduled(cron = "0 1 0 1 * ?") @@ -158,12 +164,14 @@ public class ExecutionCurrencyComponent implements IExecutionUploadComponent { LocalDate today = LocalDate.now(); log.debug("Start check new S_TRADE at {}, fromId={}", today, fromId); //выбираем STrades на сегодня с правильным section - Collection sTrades; + List sTrades; if (cleanLoad) { - sTrades = sTradeImdg.getCollectionObjectsByFieldValues(Map.of( + sTrades = CollectionUtil.castOrCopyToModifyableList( + sTradeImdg.getCollectionObjectsByFieldValues(Map.of( "tradeDate", LocalDate.now(), "section", Section.CURR.getKey() - )); + )) + ); } else { ImdgPredicateBuilder pb = sTradeImdg.predicateBuilder(); ImdgPredicate filter = pb.and( @@ -171,7 +179,9 @@ public class ExecutionCurrencyComponent implements IExecutionUploadComponent { pb.equals("section", Section.CURR.getKey()), pb.greatEqual("id", fromId) ); - sTrades = sTradeImdg.getCollectionObjectsByPredicate(filter); + sTrades = CollectionUtil.castOrCopyToModifyableList( + sTradeImdg.getCollectionObjectsByPredicate(filter) + ); } log.info("Found {} s_trade for today (fromId={})", sTrades.size(), fromId); @@ -202,14 +212,29 @@ public class ExecutionCurrencyComponent implements IExecutionUploadComponent { } log.info("{} strades left after already-added filtering", sTrades.size()); + //защита от бесконтрольно разрастающегося цикла + int iterationsWarn = iterationsWarn(sTrades); + int iterationsMax = iterationsError(sTrades); Map execsToInsert = new HashMap<>(); - for (STrades sTrd : sTrades) { + for (int i = 0; i < sTrades.size(); i++) { + if (cycleExceptionalInterrupt(i, iterationsWarn, iterationsMax)) { + return; + } + STrades sTrd = sTrades.get(i); log.trace("S_TRADE[{}] new", sTrd.getId()); //проверки IValidator validator = valFor(sTrd); Optional error = validator.tillFirstError(); + STrades counter = validator.getStored(ValidationStored.STradesCounterTrade); + if (counter != null && fromId != null && counter.getId() < fromId) { + sTrades.add(counter); + } if (error.isPresent()) { + if (error.get().getSubject().equals(ClearingError.CounterStradeNotFound)) { + log.info("STrade № {} id {} not processed. should be paired in next cycles.", sTrd.getTradeNum(), sTrd.getId()); + continue; + } String err = format("Сделка № %d не смогла быть обработана. %s", sTrd.getTradeNum(), msgResolver.resolve(error.get())); logNotifyE(err); continue; @@ -286,6 +311,8 @@ public class ExecutionCurrencyComponent implements IExecutionUploadComponent { Company company = validator.getStored(ValidationStored.STradesCompany); Company counterCompany = validator.getStored(ValidationStored.STradesCounterCompany); TradingClearingRegistry rgstr = validator.getStored(ValidationStored.STradesTradingClearingRegistry); + STrades counterSTrades = validator.getStored(ValidationStored.STradesCounterTrade); + TradingClearingRegistry counterTCR = validator.getStored(ValidationStored.STradesCounterTradingClearingRegistry); Listing listing = listingImdg.getFirstObjectByFieldValues( Map.of( "securityId", security.getId(), @@ -373,20 +400,8 @@ public class ExecutionCurrencyComponent implements IExecutionUploadComponent { eCurrency.setSecurityId(security.getId()); eCurrency.setCounterPartyId(counterCompany.getId()); { - ImdgPredicateBuilder strPb = sTradeImdg.predicateBuilder(); - STrades counterSTrades = sTradeImdg.getFirstObjectByPredicate(strPb.and( - strPb.equals("tradeNum", sTrades.getTradeNum()), - strPb.equals("section", sTrades.getSection()), - strPb.equals("classCode", sTrades.getClassCode()), - strPb.not(strPb.equals("operation", sTrades.getOperation())))); - if (counterSTrades != null) { - eCurrency.setCounterPartyTradingClearingRegistry(counterSTrades.getAccount()); - searchTcrByStrades(counterSTrades) - .ifPresent(tcr -> eCurrency.setCounterPartyTradingClearingRegistryId(tcr.getId())); - } else { - log.trace("Counter STRade for {tradeNum={}, section={}, operation={}} not found.", - sTrades.getTradeNum(), sTrades.getSection(), sTrades.getOperation()); - } + eCurrency.setCounterPartyTradingClearingRegistry(counterSTrades.getAccount()); + eCurrency.setCounterPartyTradingClearingRegistryId(counterTCR.getId()); } return eCurrency; } @@ -416,45 +431,6 @@ public class ExecutionCurrencyComponent implements IExecutionUploadComponent { } } - private void logError(IEnumId subject, Object... args) { - String err = msgResolver.resolve(new EnumMessage(subject, args)); - log.info("{}", err); -// notifications.sendNotification(ObjectType.vfrs, err, Priority.HIGH); - } - - private void logError(Long sTradeNum, EnumMessage msg) { - String err = format("Сделка № %d не смогла быть обработана. %s", sTradeNum, msgResolver.resolve(msg)); - log.warn(err); - notifications.sendNotification(ObjectType.vfrs, err, Priority.HIGH); - } - - private void logNotifyW(String msg) { - logNotify(msg, Level.WARN); - } - - private void logNotifyE(String msg) { - logNotify(msg, Level.ERROR); - } - - private void logNotify(String msg, Level lvl) { - switch (lvl) { - case WARN -> log.warn(msg); - case ERROR -> log.error(msg); - default -> log.error("lvl not supported {}. {}", lvl, msg); - } - notifications.sendNotification(ObjectType.vfrs, msg, Priority.HIGH); - } - - private Optional searchTcrByStrades(STrades sTrades) { - IValidator counterValidator = valFor(sTrades); - Optional err = counterValidator.tillFirstError(); - if (err.isPresent()) { - log.error("couldn't extract setCounterPartyTradingClearingRegistryId from strades.id={}", sTrades.getId()); - return Optional.empty(); - } - return Optional.ofNullable(counterValidator.getStored(ValidationStored.STradesTradingClearingRegistry)); - } - public IValidator valFor(STrades sTrades) { ImdgValidationContext context = new ImdgValidationContext<>(); context.setValidatedObject(sTrades); @@ -465,6 +441,7 @@ public class ExecutionCurrencyComponent implements IExecutionUploadComponent { context.addImdg(IMDGDistributedNames.Map_Company, cmpImdg); context.addImdg(IMDGDistributedNames.Map_TradingClearingRegistry, tcrImdg); context.addImdg(IMDGDistributedNames.Map_DigitalCertificateSecurity, imdgDigitCertSecurity); + context.addImdg(IMDGDistributedNames.Map_STrades, sTradeImdg); context.setLogPrefix(LogPrefixId.INSTANCE); return new ValidatorImpl<>(context, new STradesSecurityPresentValidationRule(valCash.fixIncSecCash, @@ -478,11 +455,19 @@ public class ExecutionCurrencyComponent implements IExecutionUploadComponent { new CompanyByTradingCodeCashingValidationRule<>( STrades::getCpFirmId, valCash.cmpCash, ValidationStored.STradesCounterCompany ), - new TcrByCodeAndCmpIsActiveCashingValidationRule(valCash.tcrCash), + new TcrByCodeAndCmpIsActiveCashingValidationRuleV2(valCash.tcrCash, + ImdgValidationContext::getValidatedObject, + ValidationStored.STradesCompany, + ValidationStored.STradesTradingClearingRegistry), FieldNotBlankRequiredRule.instance( "classCode", STrades::getClassCode, - ClearingError.RequiredFieldEmpty, true) + ClearingError.RequiredFieldEmpty, true), + new CounterStradePresentValidationRule(sTrades), + new TcrByCodeAndCmpIsActiveCashingValidationRuleV2(valCash.tcrCash, + ctx -> ctx.getStoredObject(ValidationStored.STradesCounterTrade), + ValidationStored.STradesCounterCompany, + ValidationStored.STradesCounterTradingClearingRegistry) ); } 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 4b6f5c1da..77e39a4e7 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 @@ -5,15 +5,14 @@ import java.time.Instant; import java.time.LocalDate; import java.time.format.DateTimeFormatter; 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 static java.lang.String.format; import org.apache.kafka.clients.producer.Producer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import org.slf4j.event.Level; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.scheduling.annotation.EnableScheduling; @@ -36,11 +35,10 @@ 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.CounterStradePresentValidationRule; import ru.spcex.clearing.service.validation.strades.STradesSecurityPresentDepositValidationRule; -import ru.spcex.clearing.service.validation.strades.TcrByCodeAndCmpIsActiveCashingValidationRule; +import ru.spcex.clearing.service.validation.strades.TcrByCodeAndCmpIsActiveCashingValidationRuleV2; import ru.spcex.platform.enumeration.MoneyFlowSide; -import ru.spcex.platform.enumeration.ObjectType; -import ru.spcex.platform.enumeration.Priority; import ru.spcex.platform.enumeration.Section; import ru.spcex.platform.enumeration.Side; import ru.spcex.platform.imdg.api.Imdg; @@ -53,8 +51,8 @@ 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.collection.CollectionUtil; import ru.spcex.platform.utils.enumeration.EnumMessage; -import ru.spcex.platform.utils.enumeration.IEnumId; import ru.spcex.platform.utils.enumeration.IEnumKey; import ru.spcex.platform.utils.enumeration.IMessageResolver; import ru.spcex.platform.utils.log.ExceptionUtils; @@ -68,7 +66,7 @@ import ru.spcex.platform.utils.validation.ValidatorImpl; */ @Component @EnableScheduling -public class ExecutionDepositComponent implements IExecutionUploadComponent { +public class ExecutionDepositComponent extends ExecutionAbstractComponent { private final Logger log = LoggerFactory.getLogger(getClass()); private final Imdg sTradeImdg; @@ -97,6 +95,7 @@ public class ExecutionDepositComponent implements IExecutionUploadComponent { NotificationSender notifications, @Qualifier("kafkaSenderWithoutRequestInfo") KafkaSender kafkaSender, IMessageResolver msgResolver) { + super(notifications, 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); @@ -111,6 +110,11 @@ public class ExecutionDepositComponent implements IExecutionUploadComponent { resetTradingDay(); } + @Override + protected Logger log() { + return log; + } + /** * Сбрасывать каждый день в 01:00:01 "0 1 0 1 * ?" */ @@ -135,12 +139,14 @@ public class ExecutionDepositComponent implements IExecutionUploadComponent { } log.debug("Start check new S_TRADE after {}, fromId={}", tradingDay, fromId); //выбираем STrades на сегодня с правильным section - Collection sTrades; + List sTrades; if (cleanLoad) { - sTrades = sTradeImdg.getCollectionObjectsByFieldValues(Map.of( + sTrades = CollectionUtil.castOrCopyToModifyableList( + sTradeImdg.getCollectionObjectsByFieldValues(Map.of( "tradeDate", LocalDate.now(), "section", Section.MKR.getKey() - )); + )) + ); } else { ImdgPredicateBuilder pb = sTradeImdg.predicateBuilder(); ImdgPredicate filter = pb.and( @@ -148,7 +154,9 @@ public class ExecutionDepositComponent implements IExecutionUploadComponent { pb.equals("section", Section.MKR.getKey()), pb.greatEqual("id", fromId) ); - sTrades = sTradeImdg.getCollectionObjectsByPredicate(filter); + sTrades = CollectionUtil.castOrCopyToModifyableList( + sTradeImdg.getCollectionObjectsByPredicate(filter) + ); } log.info("Found {} s_trade for today (fromId={})", sTrades.size(), fromId); @@ -180,14 +188,28 @@ public class ExecutionDepositComponent implements IExecutionUploadComponent { } log.info("{} strades left after already-added filtering", sTrades.size()); + //защита от бесконтрольно разрастающегося цикла + int iterationsWarn = iterationsWarn(sTrades); + int iterationsMax = iterationsError(sTrades); Map execsToInsert = new HashMap<>(); - for (STrades sTrd : sTrades) { + for (int i = 0; i < sTrades.size(); i++) { + if (cycleExceptionalInterrupt(i, iterationsWarn, iterationsMax)) { + return; + } + STrades sTrd = sTrades.get(i); log.trace("S_TRADE[{}] new", sTrd.getId()); - //проверки IValidator validator = valFor(sTrd); Optional error = validator.tillFirstError(); + STrades counter = validator.getStored(ValidationStored.STradesCounterTrade); + if (counter != null && fromId != null && counter.getId() < fromId) { + sTrades.add(counter); + } if (error.isPresent()) { + if (error.get().getSubject().equals(ClearingError.CounterStradeNotFound)) { + log.info("STrade № {} id {} not processed. should be paired in next cycles.", sTrd.getTradeNum(), sTrd.getId()); + continue; + } String err = format("Сделка № %d не смогла быть обработана. %s", sTrd.getTradeNum(), msgResolver.resolve(error.get())); logNotifyE(err); continue; @@ -253,6 +275,8 @@ public class ExecutionDepositComponent implements IExecutionUploadComponent { Company company = validator.getStored(ValidationStored.STradesCompany); Company counterCompany = validator.getStored(ValidationStored.STradesCounterCompany); TradingClearingRegistry rgstr = validator.getStored(ValidationStored.STradesTradingClearingRegistry); + STrades counterSTrades = validator.getStored(ValidationStored.STradesCounterTrade); + TradingClearingRegistry counterTCR = validator.getStored(ValidationStored.STradesCounterTradingClearingRegistry); Listing listing = listingImdg.getFirstObjectByFieldValues(Map.of("securityId", security.getId())); ExecutionDeposit eDeposit = new ExecutionDeposit(); @@ -320,17 +344,8 @@ public class ExecutionDepositComponent implements IExecutionUploadComponent { eDeposit.setCounterPartyId(counterCompany.getId()); eDeposit.setSecurityCode(sTrades.getSecCode()); { - ImdgPredicateBuilder strPb = sTradeImdg.predicateBuilder(); - STrades counterSTrades = sTradeImdg.getFirstObjectByPredicate(strPb.and( - strPb.equals("tradeNum", sTrades.getTradeNum()), - strPb.equals("section", sTrades.getSection()), - strPb.equals("classCode", sTrades.getClassCode()), - strPb.not(strPb.equals("operation", sTrades.getOperation())))); - if (counterSTrades != null) { - eDeposit.setCounterPartyTradingClearingRegistry(counterSTrades.getAccount()); - searchTcrByStrades(counterSTrades) - .ifPresent(tcr -> eDeposit.setCounterPartyTradingClearingRegistryId(tcr.getId())); - } + eDeposit.setCounterPartyTradingClearingRegistry(counterSTrades.getAccount()); + eDeposit.setCounterPartyTradingClearingRegistryId(counterTCR.getId()); } return eDeposit; } @@ -351,46 +366,6 @@ public class ExecutionDepositComponent implements IExecutionUploadComponent { } } - private void logError(IEnumId subject, Object... args) { - String err = msgResolver.resolve(new EnumMessage(subject, args)); - log.info("{}", err); -// notifications.sendNotification(ObjectType.vfrs, err, Priority.HIGH); - } - - - private void logError(Long sTradeNum, EnumMessage msg) { - String err = format("Сделка № %d не смогла быть обработана. %s", sTradeNum, msgResolver.resolve(msg)); - log.warn(err); - notifications.sendNotification(ObjectType.vfrs, err, Priority.HIGH); - } - - private void logNotifyW(String msg) { - logNotify(msg, Level.WARN); - } - - private void logNotifyE(String msg) { - logNotify(msg, Level.ERROR); - } - - private void logNotify(String msg, Level lvl) { - switch (lvl) { - case WARN -> log.warn(msg); - case ERROR -> log.error(msg); - default -> log.error("lvl not supported {}. {}", lvl, msg); - } - notifications.sendNotification(ObjectType.vfrs, msg, Priority.HIGH); - } - - private Optional searchTcrByStrades(STrades sTrades) { - IValidator counterValidator = valFor(sTrades); - Optional err = counterValidator.tillFirstError(); - if (err.isPresent()) { - log.error("couldn't extract setCounterPartyTradingClearingRegistryId from strades.id={}", sTrades.getId()); - return Optional.empty(); - } - return Optional.ofNullable(counterValidator.getStored(ValidationStored.STradesTradingClearingRegistry)); - } - @Override public void initExecCash() { ExecutionUploadCashUtil.loadExecutions(executionDepositImdg, cash); @@ -403,6 +378,7 @@ public class ExecutionDepositComponent implements IExecutionUploadComponent { context.addImdg(IMDGDistributedNames.Map_MoneyMarketSecurity, imdgMoneyMarketSecurity); context.addImdg(IMDGDistributedNames.Map_Company, imdgCompany); context.addImdg(IMDGDistributedNames.Map_TradingClearingRegistry, imdgTradingClearingRegistry); + context.addImdg(IMDGDistributedNames.Map_STrades, sTradeImdg); context.setLogPrefix(LogPrefixId.INSTANCE); return new ValidatorImpl<>(context, new STradesSecurityPresentDepositValidationRule(valCash.mmsCash), @@ -412,7 +388,15 @@ public class ExecutionDepositComponent implements IExecutionUploadComponent { new CompanyByTradingCodeCashingValidationRule<>( STrades::getCpFirmId, valCash.cmpCash, ValidationStored.STradesCounterCompany ), - new TcrByCodeAndCmpIsActiveCashingValidationRule(valCash.tcrCash) + new TcrByCodeAndCmpIsActiveCashingValidationRuleV2(valCash.tcrCash, + ImdgValidationContext::getValidatedObject, + ValidationStored.STradesCompany, + ValidationStored.STradesTradingClearingRegistry), + new CounterStradePresentValidationRule(sTrades), + new TcrByCodeAndCmpIsActiveCashingValidationRuleV2(valCash.tcrCash, + ctx -> ctx.getStoredObject(ValidationStored.STradesCounterTrade), + ValidationStored.STradesCounterCompany, + ValidationStored.STradesCounterTradingClearingRegistry) ); } 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 ce28f33c3..deed1f86a 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 @@ -5,14 +5,13 @@ import java.math.RoundingMode; import java.time.Instant; import java.time.LocalDate; 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 static java.lang.String.format; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import org.slf4j.event.Level; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.scheduling.annotation.EnableScheduling; @@ -45,9 +44,10 @@ 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.CounterStradePresentValidationRule; import ru.spcex.clearing.service.validation.strades.STradesSecurityDGCTValidationRule; import ru.spcex.clearing.service.validation.strades.STradesSecurityPresentValidationRule; -import ru.spcex.clearing.service.validation.strades.TcrByCodeAndCmpIsActiveCashingValidationRule; +import ru.spcex.clearing.service.validation.strades.TcrByCodeAndCmpIsActiveCashingValidationRuleV2; import ru.spcex.platform.enumeration.CurrencyCode; import ru.spcex.platform.enumeration.InstrumentType; import ru.spcex.platform.enumeration.ObjectType; @@ -65,8 +65,8 @@ 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.collection.CollectionUtil; import ru.spcex.platform.utils.enumeration.EnumMessage; -import ru.spcex.platform.utils.enumeration.IEnumId; import ru.spcex.platform.utils.enumeration.IEnumKey; import ru.spcex.platform.utils.enumeration.IMessageResolver; import ru.spcex.platform.utils.log.ExceptionUtils; @@ -78,7 +78,7 @@ import ru.spcex.platform.utils.validation.ValidatorImpl; @Component @EnableScheduling -public class ExecutionFondComponent implements IExecutionUploadComponent { +public class ExecutionFondComponent extends ExecutionAbstractComponent { private final Logger log = LoggerFactory.getLogger(getClass()); private final Imdg sTradeImdg; @@ -115,6 +115,7 @@ public class ExecutionFondComponent implements IExecutionUploadComponent { @Qualifier("kafkaSenderWithoutRequestInfo") KafkaSender kafkaSender, NotificationSender notifications, IMessageResolver msgResolver) { + super(notifications, msgResolver); this.sTradeImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_STrades, STrades.class); this.executionFondImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionFond, ExecutionFond.class); this.listingImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Listing, Listing.class); @@ -140,6 +141,11 @@ public class ExecutionFondComponent implements IExecutionUploadComponent { resetTradingDay(); } + @Override + protected Logger log() { + return log; + } + /** * Сбрасывать каждый день в 01:00:01 "0 1 0 1 * ?" */ @@ -164,12 +170,14 @@ public class ExecutionFondComponent implements IExecutionUploadComponent { } log.debug("Start check new S_TRADE after {}, fromId={}", tradingDay, fromId); //выбираем STrades на сегодня с правильным section - Collection sTrades; + List sTrades; if (cleanLoad) { - sTrades = sTradeImdg.getCollectionObjectsByFieldValues(Map.of( + sTrades = CollectionUtil.castOrCopyToModifyableList( + sTradeImdg.getCollectionObjectsByFieldValues(Map.of( "tradeDate", LocalDate.now(), "section", Section.FOND.getKey() - )); + )) + ); } else { ImdgPredicateBuilder pb = sTradeImdg.predicateBuilder(); ImdgPredicate filter = pb.and( @@ -177,7 +185,9 @@ public class ExecutionFondComponent implements IExecutionUploadComponent { pb.equals("section", Section.FOND.getKey()), pb.greatEqual("id", fromId) ); - sTrades = sTradeImdg.getCollectionObjectsByPredicate(filter); + sTrades = CollectionUtil.castOrCopyToModifyableList( + sTradeImdg.getCollectionObjectsByPredicate(filter) + ); } log.info("Found {} s_trade for today (fromId={})", sTrades.size(), fromId); @@ -209,14 +219,29 @@ public class ExecutionFondComponent implements IExecutionUploadComponent { } log.info("{} strades left after already-added filtering", sTrades.size()); + //защита от бесконтрольно разрастающегося цикла + int iterationsWarn = iterationsWarn(sTrades); + int iterationsMax = iterationsError(sTrades); Map execsToInsert = new HashMap<>(); - for (STrades sTrd : sTrades) { + for (int i = 0; i < sTrades.size(); i++) { + if (cycleExceptionalInterrupt(i, iterationsWarn, iterationsMax)) { + return; + } + STrades sTrd = sTrades.get(i); log.info("new S_TRADE[{}], valuation {}", sTrd.getId(), valuation); //проверки IValidator validator = valFor(sTrd); Optional error = validator.tillFirstError(); + STrades counter = validator.getStored(ValidationStored.STradesCounterTrade); + if (counter != null && fromId != null && counter.getId() < fromId) { + sTrades.add(counter); + } if (error.isPresent()) { + if (error.get().getSubject().equals(ClearingError.CounterStradeNotFound)) { + log.info("STrade № {} id {} not processed. should be paired in next cycles.", sTrd.getTradeNum(), sTrd.getId()); + continue; + } String err = format("Сделка № %d не смогла быть обработана. %s", sTrd.getTradeNum(), msgResolver.resolve(error.get())); logNotifyE(err); continue; @@ -283,6 +308,7 @@ public class ExecutionFondComponent implements IExecutionUploadComponent { Company company = validator.getStored(ValidationStored.STradesCompany); Company counterCompany = validator.getStored(ValidationStored.STradesCounterCompany); TradingClearingRegistry rgstr = validator.getStored(ValidationStored.STradesTradingClearingRegistry); + STrades counterSTrades = validator.getStored(ValidationStored.STradesCounterTrade); ClientCode clientCode = clientCodeImdg.getFirstObjectByFieldValues(Map.of("code", sTrades.getClientCode())); ExecutionFond eFond = new ExecutionFond(); @@ -413,49 +439,10 @@ public class ExecutionFondComponent implements IExecutionUploadComponent { eFond.setSettlementAmount(sTrades.getValue()); } eFond.setCounterPartyId(counterCompany.getId()); - { - ImdgPredicateBuilder strPb = sTradeImdg.predicateBuilder(); - STrades counterSTrades = sTradeImdg.getFirstObjectByPredicate(strPb.and( - strPb.equals("tradeNum", sTrades.getTradeNum()), - strPb.equals("section", sTrades.getSection()), - strPb.equals("classCode", sTrades.getClassCode()), - strPb.not(strPb.equals("operation", sTrades.getOperation())))); - if (counterSTrades != null) { - eFond.setCounterPartyTradingClearingRegistry(counterSTrades.getAccount()); - } - } + eFond.setCounterPartyTradingClearingRegistry(counterSTrades.getAccount()); return eFond; } - private void logError(IEnumId subject, Object... args) { - String err = msgResolver.resolve(new EnumMessage(subject, args)); - log.info("{}", err); -// notifications.sendNotification(ObjectType.vfrs, err, Priority.HIGH); - } - - private void logError(Long sTradeId, EnumMessage msg) { - String err = format("Сделка № %d не смогла быть обработана. %s", sTradeId, msgResolver.resolve(msg)); - log.warn(err); - notifications.sendNotification(ObjectType.vfrs, err, Priority.HIGH); - } - - private void logNotifyW(String msg) { - logNotify(msg, Level.WARN); - } - - private void logNotifyE(String msg) { - logNotify(msg, Level.ERROR); - } - - private void logNotify(String msg, Level lvl) { - switch (lvl) { - case WARN -> log.warn(msg); - case ERROR -> log.error(msg); - default -> log.error("lvl not supported {}. {}", lvl, msg); - } - notifications.sendNotification(ObjectType.vfrs, msg, Priority.HIGH); - } - @Override public void initExecCash() { ExecutionUploadCashUtil.loadExecutions(executionFondImdg, cash); @@ -473,6 +460,7 @@ public class ExecutionFondComponent implements IExecutionUploadComponent { context.addImdg(IMDGDistributedNames.Map_TradingClearingRegistry, imdgTradingClearingRegistry); context.addImdg(IMDGDistributedNames.Map_DigitalCertificateSecurity, imdgDigitCertSecurity); context.addImdg(IMDGDistributedNames.Map_UtilitarianDigitalRight, udrImdg); + context.addImdg(IMDGDistributedNames.Map_STrades, sTradeImdg); context.setLogPrefix(LogPrefixId.INSTANCE); return new ValidatorImpl<>(context, new STradesSecurityPresentValidationRule(valCash.fixIncSecCash, @@ -486,9 +474,12 @@ public class ExecutionFondComponent implements IExecutionUploadComponent { new CompanyByTradingCodeCashingValidationRule<>( STrades::getCpFirmId, valCash.cmpCash, ValidationStored.STradesCounterCompany ), - new TcrByCodeAndCmpIsActiveCashingValidationRule(valCash.tcrCash), - new STradesSecurityDGCTValidationRule(valCash.udrCash) - ); + new TcrByCodeAndCmpIsActiveCashingValidationRuleV2(valCash.tcrCash, + ImdgValidationContext::getValidatedObject, + ValidationStored.STradesCompany, + ValidationStored.STradesTradingClearingRegistry), + new STradesSecurityDGCTValidationRule(valCash.udrCash), + new CounterStradePresentValidationRule(sTrades)); } private static class ExecFondCompCash extends CashCloser { diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/ValidationStored.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/ValidationStored.java index 7ba69e2c6..012717fcb 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/ValidationStored.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/ValidationStored.java @@ -4,6 +4,8 @@ public enum ValidationStored { STradesCompany, STradesCounterCompany, STradesSecurity, STradesTradingClearingRegistry, Account, Company, STradesUtilitarianDigitalRight, + STradesCounterTrade, STradesCounterTradingClearingRegistry, + Sdf57CompanyDeb, Sdf57CompanyCred, Sdf57AccountDeb, Sdf57AccountCred, Sdf01Company, Sdf01Account, diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/strades/CounterStradePresentValidationRule.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/strades/CounterStradePresentValidationRule.java new file mode 100644 index 000000000..18bb91529 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/strades/CounterStradePresentValidationRule.java @@ -0,0 +1,44 @@ +package ru.spcex.clearing.service.validation.strades; + +import java.util.Optional; +import ru.clearing.classes.statics.data.misc.STrades; +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.ImdgPredicateBuilder; +import ru.spcex.platform.imdg.validation.ImdgValidationContext; +import ru.spcex.platform.utils.enumeration.EnumMessage; +import ru.spcex.platform.utils.validation.IValidationRule; + +public class CounterStradePresentValidationRule implements IValidationRule> { + + private final STrades sTrades; + + public CounterStradePresentValidationRule(STrades sTrades) { + this.sTrades = sTrades; + } + + @Override + public Optional validate(ImdgValidationContext context) { + Imdg sTradeImdg = context.obtainMap(IMDGDistributedNames.Map_STrades, STrades.class); + ImdgPredicateBuilder strPb = sTradeImdg.predicateBuilder(); + STrades counterSTrades = sTradeImdg.getFirstObjectByPredicate(strPb.and( + strPb.equals("tradeNum", sTrades.getTradeNum()), + strPb.equals("section", sTrades.getSection()), + strPb.equals("classCode", sTrades.getClassCode()), + strPb.not(strPb.equals("operation", sTrades.getOperation())))); + + + if (counterSTrades != null) { + context.storeObject(ValidationStored.STradesCounterTrade, counterSTrades); + return empty(); + } + return of(ClearingError.CounterStradeNotFound); + } + + @Override + public String ruleName() { + return "counter-strades-present"; + } +} 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/TcrByCodeAndCmpIsActiveCashingValidationRuleV2.java similarity index 79% rename from clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/strades/TcrByCodeAndCmpIsActiveCashingValidationRule.java rename to clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/strades/TcrByCodeAndCmpIsActiveCashingValidationRuleV2.java index 6c1048ddb..4ce96ac5f 100644 --- 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/TcrByCodeAndCmpIsActiveCashingValidationRuleV2.java @@ -1,6 +1,7 @@ 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.clearing.classes.statics.data.misc.STrades; import ru.clearing.classes.statics.data.registry.TradingClearingRegistry; @@ -18,20 +19,29 @@ 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> { +public class TcrByCodeAndCmpIsActiveCashingValidationRuleV2 implements IValidationRule> { private final CashV2ByIdAndString tcrCash; - public TcrByCodeAndCmpIsActiveCashingValidationRule(CashV2ByIdAndString tcrCash) { + public TcrByCodeAndCmpIsActiveCashingValidationRuleV2(CashV2ByIdAndString tcrCash, Function, STrades> tradeExtraction, Enum cmpEnum, Enum tcrEnum) { this.tcrCash = tcrCash; + this.tradeExtraction = tradeExtraction; + this.cmpEnum = cmpEnum; + this.tcrEnum = tcrEnum; } + private final Function, STrades> tradeExtraction; + + private final Enum cmpEnum; + + private final Enum tcrEnum; + @Override public Optional validate(ImdgValidationContext context) { - STrades validatedObject = context.getValidatedObject(); + STrades validatedObject = tradeExtraction.apply(context); if (TextUtil.isEmpty(validatedObject.getAccount())) { return of(ClearingError.TradingClearingRegistryNotFound, "", validatedObject.getAccount()); } - if (context.getStoredObject(ValidationStored.STradesCompany) == null) { //maybe unnecessary + if (context.getStoredObject(cmpEnum) == null) { //maybe unnecessary return of(ClearingError.CompanyNotFound, validatedObject.getSecCode()); } Imdg tcrImdg = context.obtainMap( diff --git a/platform-parent/platform-utils/src/main/java/ru/spcex/platform/utils/collection/CollectionUtil.java b/platform-parent/platform-utils/src/main/java/ru/spcex/platform/utils/collection/CollectionUtil.java new file mode 100644 index 000000000..124abda6d --- /dev/null +++ b/platform-parent/platform-utils/src/main/java/ru/spcex/platform/utils/collection/CollectionUtil.java @@ -0,0 +1,16 @@ +package ru.spcex.platform.utils.collection; + +import java.util.ArrayList; +import java.util.Collection; +import java.util.LinkedList; +import java.util.List; + +public class CollectionUtil { + + public static List castOrCopyToModifyableList(Collection c) { + if (c instanceof ArrayList || c instanceof LinkedList) { + return (List) c; + } + return new ArrayList<>(c); + } +}