From be93eb1ded5dae914155c77dea2ebb3c59872c2f Mon Sep 17 00:00:00 2001 From: ialbert Date: Thu, 4 May 2023 12:37:47 +0300 Subject: [PATCH] http://jira.mfd.msk:8088/browse/CLS-278 --- .../service/ExecutionDepositComponent.java | 2 + .../service/ExecutionFondComponent.java | 218 ++++++++++++++++++ .../ru/spcex/platform/enumeration/Market.java | 2 +- .../cud/registry/DealRegisterNewRequest.java | 20 +- .../domain/cud/registry/ExecutionType.java | 5 + .../json/deserialize/EnumDeserializer.java | 27 +++ .../ExecutionTypeDeserializer.java | 11 + 7 files changed, 280 insertions(+), 5 deletions(-) create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/ExecutionFondComponent.java create mode 100644 platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/registry/ExecutionType.java create mode 100644 platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/json/deserialize/EnumDeserializer.java create mode 100644 platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/json/deserialize/enumeration/ExecutionTypeDeserializer.java 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 index 5b0e2bc8e..81cddeb36 100644 --- 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 @@ -19,6 +19,7 @@ import ru.spcex.clearing.error.ClearingException; import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.platform.messaging.domain.Consts; import ru.spcex.clearing.platform.messaging.domain.cud.registry.DealRegisterNewRequest; +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.platform.enumeration.*; @@ -151,6 +152,7 @@ public class ExecutionDepositComponent { DealRegisterNewRequest requestPayload = new DealRegisterNewRequest(); requestPayload.setExecutionId(forED.getId()); requestPayload.setExchangeExecutionId(forED.getExchangeExecutionId()); + requestPayload.setExecutionType(ExecutionType.deposit); Long generatedRequestId = kafkaSender.sendRequestToQueue(destination, requestPayload); if (generatedRequestId == null) { throw new RuntimeException("failed to send request to " + destination); diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/ExecutionFondComponent.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/ExecutionFondComponent.java new file mode 100644 index 000000000..b6abd2ca6 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/ExecutionFondComponent.java @@ -0,0 +1,218 @@ +package ru.spcex.clearing.service; + +import org.apache.kafka.clients.producer.Producer; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.scheduling.annotation.EnableScheduling; +import org.springframework.scheduling.annotation.Scheduled; +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.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.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.Consts; +import ru.spcex.clearing.platform.messaging.domain.cud.registry.DealRegisterNewRequest; +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.platform.enumeration.*; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; +import ru.spcex.platform.utils.enumeration.*; +import ru.spcex.platform.utils.log.ExceptionUtils; +import ru.spcex.platform.utils.time.TimeUtil; +import ru.spcex.platform.utils.validation.IValidator; + +import java.time.Instant; +import java.time.LocalDate; +import java.util.Collection; +import java.util.Map; +import java.util.Optional; +import java.util.function.Function; + +@Component +@EnableScheduling +public class ExecutionFondComponent { + private final Logger log = LoggerFactory.getLogger(getClass()); + + private final Imdg sTradeImdg; + private final Imdg executionFondImdg; + private final Imdg listingImdg; + private final Imdg clientCodeImdg; + + //fixme ждать ТЗ + Long tradeNum; + Instant tradingDay; + private final IMessageResolver msgResolver = new SimpleMessageResolver(); + private final Function stradesValidator; + private final KafkaSender kafkaSender; + + + @Autowired + public ExecutionFondComponent(ImdgProvider imdgProvider, Producer kafka, + @Qualifier("sTradesValidator") Function stradesValidator, + @Qualifier("kafkaSenderWithoutRequestInfo") KafkaSender kafkaSender) { + 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); + this.clientCodeImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ClientCode, ClientCode.class); + this.stradesValidator = stradesValidator; + this.kafkaSender = kafkaSender; + + resetTradingDay(); + } + + /** + * Сбрасывать каждый день в 01:00:01 "0 1 0 1 * ?" + */ + @Scheduled(cron = "0 1 0 1 * ?") + public void resetTradingDay() { + log.trace("Recheck today trading day for search STrade. Current state: tradeNum={}, tradeDat={}", tradeNum, tradingDay); + 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() { + log.debug("Start check new S_TRADE after {}", tradingDay); + //выбираем STrades на сегодня с правильным section + Collection sTrades = sTradeImdg.getCollectionObjectsByFieldValues(Map.of( + "tradeDate", LocalDate.now(), + "section", Section.FOND.getKey() + )); + log.info("Found {} s_trade for today", sTrades.size()); + + if (sTrades.isEmpty()) { + logError(ClearingError.NewDealsNotFound); + return; + } + + //убираем уже добавленные в ExecutionDeposit + 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; + return executionFondImdg.getSingleObjectByFieldValues( + Map.of("tradingDate", sTrd.getTradeDate(), + "exchangeExecutionId", sTrd.getTradeNum(), + "side", excDepSide.getKey())) != null; + }); + log.info("{} strades left after already-added filtering", sTrades.size()); + + 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; + } + ExecutionFond newED; + try { + newED = createExecutionFond(sTrd, validator); + executionFondImdg.insert(newED); + sendNotification(newED); + log.debug("New executionDeposit.id={} was created.", newED.getId()); + } catch (ClearingException ce) { + auditMessage(ce); + } catch (Exception e) { + log.error("When create new ExecutionDeposit by STrade[{}] error: {}", sTrd.getId(), ExceptionUtils.getStackTrace(e)); + } + } + + Long newMaxTradeNum = sTrades.stream().mapToLong(STrades::getTradeNum).max().orElseGet(() -> tradeNum); + log.info("Process completed. Next tradeNum is {}", newMaxTradeNum); + } + + + protected void auditMessage(ClearingException ce) { + log.error("AUDIT error code {}: {}", ce.getEnumMsg(), ce.getMessage()); + } + + protected void sendNotification(ExecutionFond forEF) throws ClearingException { + final String destination = Consts.REGISTRY_DEAL_REGISTER_NEW; + DealRegisterNewRequest requestPayload = new DealRegisterNewRequest(); + requestPayload.setExecutionId(forEF.getId()); + requestPayload.setExchangeExecutionId(forEF.getExchangeExecutionId()); + requestPayload.setExecutionType(ExecutionType.fond); + Long generatedRequestId = kafkaSender.sendRequestToQueue(destination, requestPayload); + if (generatedRequestId == null) { + throw new RuntimeException("failed to send request to " + destination); + } + } + + protected ExecutionFond createExecutionFond(STrades sTrades, + IValidator validator) throws ClearingException { + Security security = validator.getStored(ValidationStored.STradesSecurity); + Company company = validator.getStored(ValidationStored.STradesCompany); + Company counterCompany = validator.getStored(ValidationStored.STradesCounterCompany); + TradingClearingRegistry rgstr = validator.getStored(ValidationStored.STradesTradingClearingRegistry); + ClientCode clientCode = clientCodeImdg.getSingleObjectByFieldValues(Map.of("code", sTrades.getClientCode())); + + ExecutionFond eFond = new ExecutionFond(); + final Instant now = Instant.now(); + eFond.setCreated(now); + eFond.setTradingDate(sTrades.getTradeDate()); + eFond.setClearingDate(TimeUtil.toLocalDate(now)); + eFond.setExchangeExecutionId(sTrades.getTradeNum()); + eFond.setExchangeExecutionTime(sTrades.getTradeDateTime()); + eFond.setTradingClearingRegistryId(rgstr.getId()); + eFond.setMarket(Market.fond.getKey()); + eFond.setPrice(sTrades.getPrice()); + eFond.setLots(sTrades.getQty()); + eFond.setQuantity(sTrades.getQtyPcs()); + { + Side sTradeSide = IEnumKey.getEnumByKey(Side.class, sTrades.getOperation()); + if (sTradeSide != null) { + switch (sTradeSide) { + case BUY -> eFond.setSide(MoneyFlowSide.BUY.getKey()); + case SELL -> eFond.setSide(MoneyFlowSide.SELL.getKey()); + } + } + } + //fixme символьный код по справочнику currencyCode, соответствующий значению из sTrades.settleCurrency + eFond.setSettlementCurrency(CurrencyCode.RUB.getKey()); + //fixme micro vs milli + eFond.setExchangeExecutionMicroseconds(Instant.ofEpochMilli(sTrades.getTradeTimeMs())); + eFond.setCompanyId(company.getId()); + eFond.setDuration(sTrades.getRepoTerm()); + eFond.setComment(sTrades.getBrokerRef()); + if (clientCode != null) { + eFond.setClientCodeId(clientCode.getId()); + } + eFond.setSettlementDate(sTrades.getSettleDate()); + eFond.setSettlementCode(sTrades.getSettleCode()); + eFond.setSecurityFullName(security.getFullName()); + eFond.setSecuritySymbol(security.getSecuritySymbol()); + eFond.setSecurityId(security.getId()); + eFond.setInterestAmount(sTrades.getAccruedint()); + eFond.setExchangeOrderId(sTrades.getOrderNum()); + eFond.setSettlementAmount(sTrades.getValue()); + eFond.setCounterPartyId(counterCompany.getId()); + return eFond; + } + + private void logError(IEnumId subject, Object... args) { + log.warn("{}", msgResolver.resolve(new EnumMessage(subject, args))); + } + + private void logError(Long sTradeId, EnumMessage msg) { + log.warn("sTrade id={} {}", sTradeId, msgResolver.resolve(msg)); + } +} diff --git a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/Market.java b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/Market.java index 8d20c67ef..e8f048253 100644 --- a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/Market.java +++ b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/Market.java @@ -3,7 +3,7 @@ package ru.spcex.platform.enumeration; import ru.spcex.platform.utils.enumeration.IEnumKey; public enum Market implements IEnumKey { - mkrs("MKRS"); + mkrs("MKRS"), fond("FOND"); private final String key; diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/registry/DealRegisterNewRequest.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/registry/DealRegisterNewRequest.java index 0e37c8711..0954ca33f 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/registry/DealRegisterNewRequest.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/registry/DealRegisterNewRequest.java @@ -2,10 +2,10 @@ package ru.spcex.clearing.platform.messaging.domain.cud.registry; import com.fasterxml.jackson.annotation.JsonProperty; - -import java.math.BigDecimal; -import java.time.Instant; -import java.time.LocalDate; +import com.fasterxml.jackson.databind.annotation.JsonDeserialize; +import com.fasterxml.jackson.databind.annotation.JsonSerialize; +import ru.spcex.clearing.platform.messaging.domain.json.deserialize.enumeration.ExecutionTypeDeserializer; +import ru.spcex.clearing.platform.messaging.domain.json.serialize.EnumSerializer; public class DealRegisterNewRequest { @@ -15,6 +15,11 @@ public class DealRegisterNewRequest { @JsonProperty public Long exchangeExecutionId; + @JsonSerialize(using = EnumSerializer.class) + @JsonDeserialize(using = ExecutionTypeDeserializer.class) + @JsonProperty + public ExecutionType executionType; + public Long getExecutionId() { return executionId; @@ -32,4 +37,11 @@ public class DealRegisterNewRequest { this.exchangeExecutionId = exchangeExecutionId; } + public ExecutionType getExecutionType() { + return executionType; + } + + public void setExecutionType(ExecutionType executionType) { + this.executionType = executionType; + } } diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/registry/ExecutionType.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/registry/ExecutionType.java new file mode 100644 index 000000000..aa6206e1e --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/registry/ExecutionType.java @@ -0,0 +1,5 @@ +package ru.spcex.clearing.platform.messaging.domain.cud.registry; + +public enum ExecutionType { + deposit, fond; +} diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/json/deserialize/EnumDeserializer.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/json/deserialize/EnumDeserializer.java new file mode 100644 index 000000000..4be4955f2 --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/json/deserialize/EnumDeserializer.java @@ -0,0 +1,27 @@ +package ru.spcex.clearing.platform.messaging.domain.json.deserialize; + +import com.fasterxml.jackson.core.JsonParseException; +import com.fasterxml.jackson.core.JsonParser; +import com.fasterxml.jackson.databind.DeserializationContext; +import com.fasterxml.jackson.databind.JsonDeserializer; + +import java.io.IOException; + +public abstract class EnumDeserializer> extends JsonDeserializer { + + @Override + public T deserialize(JsonParser p, DeserializationContext ctxt) throws IOException { + String v = p.getText(); + if (v == null || v.length() == 0) { + return null; + } + for (T enumConstant : getEnumClass().getEnumConstants()) { + if (enumConstant.name().equalsIgnoreCase(v)) { + return enumConstant; + } + } + throw new JsonParseException(p, "couldn't parse " + getEnumClass().getSimpleName() + " value '" + v + "'"); + } + + protected abstract Class getEnumClass(); +} diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/json/deserialize/enumeration/ExecutionTypeDeserializer.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/json/deserialize/enumeration/ExecutionTypeDeserializer.java new file mode 100644 index 000000000..56c50fae5 --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/json/deserialize/enumeration/ExecutionTypeDeserializer.java @@ -0,0 +1,11 @@ +package ru.spcex.clearing.platform.messaging.domain.json.deserialize.enumeration; + +import ru.spcex.clearing.platform.messaging.domain.cud.registry.ExecutionType; +import ru.spcex.clearing.platform.messaging.domain.json.deserialize.EnumDeserializer; + +public class ExecutionTypeDeserializer extends EnumDeserializer { + @Override + protected Class getEnumClass() { + return ExecutionType.class; + } +}