This commit is contained in:
parent
ac3e325269
commit
be93eb1ded
7 changed files with 280 additions and 5 deletions
|
|
@ -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);
|
||||
|
|
|
|||
|
|
@ -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<STrades> sTradeImdg;
|
||||
private final Imdg<ExecutionFond> executionFondImdg;
|
||||
private final Imdg<Listing> listingImdg;
|
||||
private final Imdg<ClientCode> clientCodeImdg;
|
||||
|
||||
//fixme ждать ТЗ
|
||||
Long tradeNum;
|
||||
Instant tradingDay;
|
||||
private final IMessageResolver msgResolver = new SimpleMessageResolver();
|
||||
private final Function<STrades, IValidator> stradesValidator;
|
||||
private final KafkaSender kafkaSender;
|
||||
|
||||
|
||||
@Autowired
|
||||
public ExecutionFondComponent(ImdgProvider imdgProvider, Producer<String, Object> kafka,
|
||||
@Qualifier("sTradesValidator") Function<STrades, IValidator> 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> 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<EnumMessage> 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));
|
||||
}
|
||||
}
|
||||
|
|
@ -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;
|
||||
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -0,0 +1,5 @@
|
|||
package ru.spcex.clearing.platform.messaging.domain.cud.registry;
|
||||
|
||||
public enum ExecutionType {
|
||||
deposit, fond;
|
||||
}
|
||||
|
|
@ -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<T extends Enum<T>> extends JsonDeserializer<T> {
|
||||
|
||||
@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<T> getEnumClass();
|
||||
}
|
||||
|
|
@ -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<ExecutionType> {
|
||||
@Override
|
||||
protected Class<ExecutionType> getEnumClass() {
|
||||
return ExecutionType.class;
|
||||
}
|
||||
}
|
||||
Loading…
Add table
Reference in a new issue