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 dfaacf56d..2e9d0341c 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 @@ -3,8 +3,11 @@ package ru.spcex.clearing.error; import ru.spcex.platform.utils.enumeration.IEnumId; public enum ClearingError implements IEnumId { + GeneralError(5400L), + RecordNotFound(5406L), CompanyCreditCheck(5412L), CompanyDebitCheck(5413L), + CompanyNotFound(5410L), ; private final Long id; diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/error/ClearingException.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/error/ClearingException.java new file mode 100644 index 000000000..e70bf1f67 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/error/ClearingException.java @@ -0,0 +1,34 @@ +package ru.spcex.clearing.error; + +import ru.spcex.platform.utils.enumeration.EnumMessage; + +public class ClearingException extends Exception { + final EnumMessage enumMsg; + + public ClearingException(EnumMessage msg) { + this.enumMsg = msg; + } + + public ClearingException(ClearingError code) { + this.enumMsg = new EnumMessage(code); + } + + public ClearingException(ClearingError code, String message) { + super(code == null ? message : code.getId() + " " + message); + this.enumMsg = new EnumMessage(code); + } + + public ClearingException(ClearingError code, String message, Throwable cause) { + super(code == null ? message : code.getId() + " " + message, cause); + this.enumMsg = new EnumMessage(code); + } + + public ClearingException(String message, Throwable cause) { + super(message, cause); + this.enumMsg = new EnumMessage(ClearingError.GeneralError); + } + + public EnumMessage getEnumMsg() { + return enumMsg; + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/ExecutionDepositComponent.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/ExecutionDepositComponent.java new file mode 100644 index 000000000..e16d34736 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/ExecutionDepositComponent.java @@ -0,0 +1,350 @@ +package ru.spcex.clearing.service; + +import com.fasterxml.jackson.annotation.JsonProperty; +import org.apache.kafka.clients.producer.Producer; +import org.apache.kafka.clients.producer.ProducerRecord; +import org.apache.kafka.clients.producer.RecordMetadata; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.scheduling.annotation.EnableScheduling; +import org.springframework.scheduling.annotation.Scheduled; +import org.springframework.stereotype.Component; +import org.springframework.util.StringUtils; +import ru.clearing.classes.statics.data.account.Account; +import ru.clearing.classes.statics.data.account.AccountBalance; +import ru.clearing.classes.statics.data.clearing.VerificationResult; +import ru.clearing.classes.statics.data.company.Company; +import ru.clearing.classes.statics.data.execution.ExecutionDeposit; +import ru.clearing.classes.statics.data.misc.Listing; +import ru.clearing.classes.statics.data.misc.STrade; +import ru.clearing.classes.statics.data.sdf.SDf01; +import ru.clearing.classes.statics.data.security.Security; +import ru.spcex.clearing.error.ClearingError; +import ru.spcex.clearing.error.ClearingException; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.platform.messaging.domain.ActionType; +import ru.spcex.clearing.platform.messaging.domain.BaseRequest; +import ru.spcex.clearing.platform.messaging.domain.Consts; +import ru.spcex.clearing.platform.messaging.domain.cud.securitites.MoneyMarketSecurityNewRequest; +import ru.spcex.platform.classes.base.interfaces.WithId; +import ru.spcex.platform.enumeration.AccountType; +import ru.spcex.platform.enumeration.Allowed; +import ru.spcex.platform.enumeration.Market; +import ru.spcex.platform.enumeration.ResultStatuses; +import ru.spcex.platform.imdg.api.Imdg; +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.utils.enumeration.EnumMessage; +import ru.spcex.platform.utils.log.ExceptionUtils; +import ru.spcex.platform.utils.time.TimeUtil; +import ru.spcex.clearing.platform.messaging.domain.cud.registry.CoveredDealRegisterNewRequest; + + +import java.math.BigDecimal; +import java.math.RoundingMode; +import java.time.Instant; +import java.time.LocalDate; +import java.util.*; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.Future; +import java.util.function.Function; +import java.util.stream.Collectors; + +/** + * 1.35. executionDeposit - Сделки + * I - Изменение executionDeposit при получении новых сделок из ТС (s_trade) + */ +@Component +@EnableScheduling +public class ExecutionDepositComponent { + private final Logger log = LoggerFactory.getLogger(getClass()); + + private final ImdgProvider imdgProvider; + private Imdg sTradeImdg; + private Imdg securityImdg; + private Imdg executionDepositImdg; + private Imdg companyImdg; + private Imdg accountImdg; + private Imdg listingImdg; + + private ImdgId idGenerator; + Producer kafka; + + Long tradeNum; + Instant tradingDay; + + + @Autowired + public ExecutionDepositComponent(ImdgProvider imdgProvider, Producer kafka) { + this.imdgProvider = imdgProvider; + this.sTradeImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_STrade, STrade.class); + this.securityImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Security, Security.class); + this.executionDepositImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionDeposit, ExecutionDeposit.class); + this.companyImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Company, Company.class); + this.accountImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Account, Account.class); + this.listingImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Listing, Listing.class); + + this.idGenerator = imdgProvider.getImdgIdGenerator(); + this.kafka = kafka; + + resetTradingDay(); + } + + /** + * Сбрасывать каждый день в 01:00:01 "0 1 0 1 * ?" + */ + @Scheduled(cron = "${clearing-service.scheduler.check-s-trade}") + public void resetTradingDay() { + Instant today = TimeUtil.localDateToInstant(LocalDate.now()); + if (tradingDay == null || !tradingDay.equals(today)) { + tradeNum = -1L; + tradingDay = today; + } + log.info("Reset trading day for search STrade: tradeNum={}, tradeDat={}", tradeNum, tradingDay); + } + + public void processNewTS() { + Long tradeNum = -1L; // todo уточнить как он обновляется + ImdgPredicateBuilder pb = sTradeImdg.predicateBuilder(); + ImdgPredicate sql = pb.and(pb.greater("tradeNum",tradeNum), pb.greatEqual("tradeDateTime", tradingDay)); + Collection sTrades = sTradeImdg.getCollectionObjectsByPredicate(sql); + log.info("Found {} new s_trade with trade_num>{}", sTrades.size(), tradeNum); + + if (sTrades.isEmpty()) { + log.info("No new sTrades."); + return; + } + + // Выявление новых сделок необходимо выполнить следующие контрольные проверки: + + // Проверить все инструменты. + { + Set secCodesOfSTrade = sTrades.stream().map(STrade::getSecCode).filter(Objects::nonNull).collect(Collectors.toSet()); + log.debug("Verify {} instruments for {} STrade's.", secCodesOfSTrade.size(), sTrades.size()); + ImdgPredicate allIn = securityImdg.predicateBuilder().in("securitySymbol", secCodesOfSTrade.toArray(new String[0])); + Collection foundSecurities = securityImdg.getCollectionObjectsByPredicate(allIn); + Set secCodesOfSecurity = foundSecurities.stream().map(Security::getSecuritySymbol).filter(Objects::nonNull).collect(Collectors.toSet()); + if (secCodesOfSecurity.containsAll(secCodesOfSTrade)) { + log.debug("All {} Security found by {} secCodes from STrade", + secCodesOfSecurity.size(), secCodesOfSTrade.size()); + } else { + HashSet notFoundSymbol = new HashSet<>(secCodesOfSTrade); + notFoundSymbol.removeAll(secCodesOfSecurity); + log.info("Found only {} Security by {} secCodes from STrade. Not found: {}", + secCodesOfSecurity.size(), secCodesOfSTrade.size(), notFoundSymbol); + createNewSecurities(notFoundSymbol); + log.info("Stop till they all will be created"); + + auditMessage("В security нет записей с securitySymbol", notFoundSymbol); + return; + } + } + + Long generationId = idGenerator.nextId(); + log.info("generationId = {}", generationId); + for (STrade trade:sTrades) { + log.trace("Check s_trade[{}].tradeNum={}", trade.getId(), trade.getTradeNum()); + Collection existsEDeposit = executionDepositImdg.getCollectionObjectsByFieldValues(Map.of( + "exchangeExecutionId", trade.getTradeNum(), + "exchangeExecutionTime", trade.getTradeDateTime() + )); + + if (existsEDeposit.isEmpty()) { + log.trace("S_TRADE[{}] new", trade.getId()); + ExecutionDeposit newED = null; + try { + newED = createExecutionDeposit(trade, Allowed.ALLOWED/*todo уточнить момент заполнения*/, generationId); + verification(newED); + executionDepositImdg.insert(newED); + sendNotification(newED); + } catch (ClearingException ce) { + auditMessage(ce); + } catch (Exception e) { + if (newED != null) { + newED.setCoverageStatus(Allowed.DENIED.getKey()); + } + log.error("When create new ExecutionDeposit by STrade[{}]", trade.getId()); + } + + } else { + long[] idToLong = existsEDeposit.stream().mapToLong(ed-> ed.getId()).toArray(); + log.warn("S_TRADE[{}] already has executionDeposit: {}", trade.getId(), Arrays.toString(idToLong)); + } + } + + + Long newMaxTradeNum = sTrades.stream().mapToLong(STrade::getTradeNum).max().orElseGet(()-> tradeNum); + log.debug("Next tradeNum is {}", newMaxTradeNum); + } + + protected void verification(ExecutionDeposit forED) throws ClearingException { + /*todo Рассчитанные в КС контрольные суммы (общее количество сделок и суммарный объем заключенных сделок в денежном выражении) + должны совпадать со значениями, рассчитанными Торговой системой: + count(execution[tradingDay]) = count (trade_arqua) + */ + // использовать ли VerificationResultComponent для сверки или здесь код добавить. + + } + + protected void auditMessage(ClearingException ce) { + log.error("AUDIT error code {}: {}", ce.getEnumMsg(), ce.getMessage()); + } + + protected void auditMessage(String message, Object... ids) { + String txt = message; + if (ids != null && ids.length > 0) { + txt += " object:" + Arrays.toString(ids); + } + log.error("audit \"clearing-service\", errorText: {}", txt); + } + + /** + * в очередь kafka для модуля securities-service сообщение о добавлении инструмента с параметром securitySymbol=s_trade.sec_code + * @param newSymbolRequest + */ + protected void createNewSecurities(Collection newSymbolRequest) { + final String destination = Consts.DESTINATION_MONEY_MARKET_SECURITY_NEW; + List symbolRequests = new ArrayList<>(newSymbolRequest); // чтобы в случае ошибки отобразить номер в логе + List> sendAll = new ArrayList<>(symbolRequests.size()); + for (String newSymbol: symbolRequests) { + if (newSymbol == null || newSymbol.isEmpty()) { + log.warn("Empty SecuritySumbol"); + } else { + MoneyMarketSecurityNewRequest requestPayload = new MoneyMarketSecurityNewRequest(); + requestPayload.setSecuritySymbol(newSymbol); + + BaseRequest request = new BaseRequest<>(); + request.setId(idGenerator.nextId()); + request.setActionType(ActionType.NEW); + request.setRequestPayload(requestPayload); + +// saveRequestToStorage(destination, request); //сохраняет данные о запросе в хранилище + log.trace("Send to {} new symbol \"{}\" ", destination, newSymbol); + Future send = kafka.send(new ProducerRecord<>(destination, request)); + sendAll.add(send); + } + } + int i = 0; + for (Future future: sendAll) { + try { + future.get(); // get exception + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + String about = i < symbolRequests.size() ? symbolRequests.get(i) : "(out of range i=" + i + ")"; + log.warn("Thread interrupted! On send symbol \"{}\"", about); + throw new RuntimeException(e); + } catch (ExecutionException e) { + String about = i < symbolRequests.size() ? symbolRequests.get(i) : "(out of range i=" + i + ")"; + log.error("Error send message for symbol \"{}\" to {}: {}", destination, + about, ExceptionUtils.getStackTrace(e.getCause() == null ? e : e.getCause())); + } + i++; + } + } + + protected void sendNotification(ExecutionDeposit forED) { + final String destination = Consts.REGISTRY_COVERED_DEAL_REGISTER_NEW; + CoveredDealRegisterNewRequest requestPayload = new CoveredDealRegisterNewRequest(); + requestPayload.setExecutionId(forED.getId()); +// requestPayload.setCompanyFullName(forED.getCompanyFullName()); + requestPayload.setTradingDate(forED.getTradingDate()); + requestPayload.setExchangeExecutionId(forED.getExchangeExecutionId()); + requestPayload.setExchangeExecutionTime(forED.getExchangeExecutionTime()); + requestPayload.setSecuritySymbol(forED.getSecuritySymbol()); + requestPayload.setSecurityFullName(forED.getSecurityFullName()); +// requestPayload.setSellerFullName(forED.getSellerFullName()); +// requestPayload.setSellerClearingCode(forED.getSellerClearingCode()); +// String requestPayload.setSellerAccount(forED.getAccountId()); +// String requestPayload.setBuyerFullName(forED.getBuyerFullName()); +// requestPayload.setBuyerClearingCode(forED.getBuyerClearingCode()); +// String requestPayload.setBuyerAccount(forED.getBuyerAccount()); +// BigDecimal requestPayload.setAmount(forED.getAmount()); + requestPayload.setId(idGenerator.nextId()); + requestPayload.setCreatedAt(forED.getCreated()); + requestPayload.setUpdatedAt(forED.getUpdated()); + requestPayload.setClearingDate(forED.getClearingDate()); + + BaseRequest request = new BaseRequest<>(); + request.setId(idGenerator.nextId()); + request.setActionType(ActionType.NEW); + request.setRequestPayload(requestPayload); + log.trace("Send to {} new ExecutionDeposit[{}]", destination, forED.getId()); + Future send = kafka.send(new ProducerRecord<>(destination, request)); + try { + send.get(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new RuntimeException(e); + } catch (ExecutionException e) { + throw new RuntimeException(e); + } + } + + protected ExecutionDeposit createExecutionDeposit(STrade sTrade, + Allowed coverageStatus, Long sessionId) throws ClearingException { + Account account = accountImdg.getSingleObjectByFieldValues(Map.of("account", sTrade.getMoneyAccount())); + Security security = securityImdg.getSingleObjectByFieldValues(Map.of("securitySymbol", sTrade.getSecCode())); + if (security == null) { + log.warn("security securitySymbol=\"{}\" not found", sTrade.getSecCode()); + throw new ClearingException(new EnumMessage(ClearingError.RecordNotFound, sTrade.getSecCode())); + } + Listing listing = null; + if (security != null) { + listing = listingImdg.getSingleObjectByFieldValues(Map.of("securityId", security.getId())); + } + Company company = companyImdg.getSingleObjectByFieldValues(Map.of("tradingCode", sTrade.getFirmId())); + if (company == null) { + throw new ClearingException(new EnumMessage(ClearingError.CompanyNotFound, sTrade.getFirmId())); + } + + return createExecutionDeposit(sTrade, account, listing, company, security, coverageStatus, sessionId); + } + + private ExecutionDeposit createExecutionDeposit(STrade sTrade, Account account, Listing listing, + Company company, Security security, + Allowed coverageStatus, Long sessionId) { + ExecutionDeposit eDeposit = new ExecutionDeposit(); + eDeposit.setId(idGenerator.nextId()); + final Instant now = Instant.now(); + final LocalDate nowDay = TimeUtil.toLocalDate(now); + eDeposit.setCreated(now); + eDeposit.setTradingDate(nowDay); + eDeposit.setClearingDate(nowDay); + + eDeposit.setExchangeExecutionId(sTrade.getTradeNum()); + eDeposit.setExchangeExecutionTime(sTrade.getTradeDateTime()); + if (account != null) { + eDeposit.setAccountId(account.getId()); + } + eDeposit.setMarket(Market.mkrs.getKey()); + eDeposit.setPrice(sTrade.getPrice()); + eDeposit.setLots(sTrade.getQty()); + if (listing != null && listing.getLotSize() != null && eDeposit.getLots() != null) { + BigDecimal quantity = eDeposit.getLots().multiply(listing.getLotSize()); + eDeposit.setQuantity(quantity); + } + eDeposit.setFirstLegAmount(sTrade.getValue()); + eDeposit.setSecondLegAmount(sTrade.getValue()); + //eDeposit.setInterestAmount(null); + eDeposit.setSide(sTrade.getOperation());//Символьный код по справочнику moneyFlowSide), соответствующий значению из s_trade.operation (sTrade.getOperation()) + eDeposit.setSettlementCurrency("RUB"); // (справочник currencyCode) + eDeposit.setCompanyId(company.getId()); + eDeposit.setDuration(sTrade.getDaysToMatDate()); + eDeposit.setFirstLegSettlementDate(nowDay); + eDeposit.setSecondLegSettlementDate(sTrade.getSettleDate()); + //eDeposit.setFirstLegSettlementCode(null); + //eDeposit.setSecondLegSettlementCode(null); + eDeposit.setSecurityFullName(security.getFullName()); + eDeposit.setSecuritySymbol(security.getSecuritySymbol()); + eDeposit.setSecurityId(security.getId()); + //eDeposit.setCounterPartyId(null); + eDeposit.setCoverageStatus(coverageStatus == null ? null : coverageStatus.getKey()); // Заполняется по справочнику allowed в результате расчета требований и обязательств. TODO + eDeposit.setSessionId(sessionId); + + return eDeposit; + } + +} diff --git a/clearing-parent/clearing-service/src/main/resources/application.properties b/clearing-parent/clearing-service/src/main/resources/application.properties index af4ebcf34..e0908f94a 100644 --- a/clearing-parent/clearing-service/src/main/resources/application.properties +++ b/clearing-parent/clearing-service/src/main/resources/application.properties @@ -19,4 +19,5 @@ clearing-service.kafka-producer.batch-size=16384 clearing-service.kafka-producer.linger-ms=1 clearing-service.kafka-producer.buffer-memory=33554432 -clearing-service.scheduler.check-payment-instruction=*/5 * * * * * \ No newline at end of file +clearing-service.scheduler.check-payment-instruction=*/5 * * * * * +clearing-service.scheduler.check-s-trade=0 1 0 1 * ?