clearing-service http://jira.mfd.msk:8088/browse/CLS-628 обработка ExecutionCurrency (доделать некоторые поля)

This commit is contained in:
AKurakin 2024-04-12 18:14:05 +03:00
parent 4580880fe4
commit 0dfd89d294
8 changed files with 443 additions and 3 deletions

View file

@ -201,6 +201,26 @@ public class ValidationConfig {
};
}
@Bean("sTradesValidatorCurrency")
public Function<STrades, IValidator> sTradesValidatorCurrency() {
return sTrades -> {
ImdgValidationContext<STrades> context = new ImdgValidationContext<>();
context.setValidatedObject(sTrades);
context.addImdg(IMDGDistributedNames.Map_CurrencyPairSecurity, currencyPairSecurityImdg);
context.addImdg(IMDGDistributedNames.Map_MoneyMarketSecurity, imdgMoneyMarketSecurity);
context.addImdg(IMDGDistributedNames.Map_FixedIncomeSecurity, fixedIncomeSecurityImdg);
context.addImdg(IMDGDistributedNames.Map_EquitySecurity, equitySecurityImdg);
context.addImdg(IMDGDistributedNames.Map_Company, imdgCompany);
context.addImdg(IMDGDistributedNames.Map_TradingClearingRegistry, imdgTradingClearingRegistry);
context.setLogPrefix(LogPrefixId.INSTANCE);
return new ValidatorImpl<>(context,
STradesValidationRule.SecurityPresentCurrencyPair,
STradesValidationRule.CompanyPresent,
STradesValidationRule.CounterCompanyPresent,
STradesValidationRule.TradingClearingRegistryPresent);
};
}
@Bean("obligationAndRequirementsAdmissionValidator")
public Function<Registry, IValidator> registryValidator() {
return rgs -> {

View file

@ -8,6 +8,7 @@ import org.springframework.scheduling.annotation.EnableScheduling;
import org.springframework.stereotype.Service;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonIdRequest;
import ru.spcex.clearing.service.execution.ExecutionCurrencyComponent;
import ru.spcex.clearing.service.execution.ExecutionDepositComponent;
import ru.spcex.clearing.service.execution.ExecutionFondComponent;
import ru.spcex.platform.utils.log.ExceptionUtils;
@ -26,18 +27,21 @@ public class ClearingService implements DisposableBean {
private final Clearing clearing;
private final ExecutionDepositComponent executionDepositComponent;
private final ExecutionFondComponent executionFondComponent;
private final ExecutionCurrencyComponent executionCurrencyComponent;
@Autowired
public ClearingService(SdfCreatorBySTLDPayment sdfCreator, PaymentUpdateBySdf04 paymentUpdater,
VerificationResultComponent verificationResultComponent, Clearing clearing,
ExecutionDepositComponent executionDepositComponent,
ExecutionFondComponent executionFondComponent) {
ExecutionFondComponent executionFondComponent,
ExecutionCurrencyComponent executionCurrencyComponent) {
this.sdfCreator = sdfCreator;
this.paymentUpdater = paymentUpdater;
this.verificationResultComponent = verificationResultComponent;
this.executionDepositComponent = executionDepositComponent;
this.clearing = clearing;
this.executionFondComponent = executionFondComponent;
this.executionCurrencyComponent = executionCurrencyComponent;
this.executor = Executors.newSingleThreadExecutor();
}
@ -106,6 +110,7 @@ public class ClearingService implements DisposableBean {
try {
executionDepositComponent.processNewTS();
executionFondComponent.processNewTS();
executionCurrencyComponent.processNewTS();
} catch (Throwable e) {
log.error("{}", ExceptionUtils.getStackTrace(e));
}

View file

@ -0,0 +1,287 @@
package ru.spcex.clearing.service.execution;
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.company.Company;
import ru.clearing.classes.statics.data.execution.ExecutionCurrency;
import ru.clearing.classes.statics.data.misc.Listing;
import ru.clearing.classes.statics.data.misc.Market;
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.Section;
import ru.spcex.platform.enumeration.Side;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
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;
import ru.spcex.platform.utils.time.TimeUtil;
import ru.spcex.platform.utils.validation.IValidator;
import java.math.BigDecimal;
import java.time.Instant;
import java.time.LocalDate;
import java.time.format.DateTimeFormatter;
import java.util.Collection;
import java.util.Map;
import java.util.Optional;
import java.util.function.Function;
/**
* ExecutionCurrency - Сделки на Валютной секции
*/
@Component
@EnableScheduling
public class ExecutionCurrencyComponent {
private final Logger log = LoggerFactory.getLogger(getClass());
private final Imdg<STrades> sTradeImdg;
private final Imdg<ExecutionCurrency> executionCurrencyImdg;
private final Imdg<Listing> listingImdg;
// protected transient Long tradeNum; //todo used?
// protected transient Instant tradingDay;
private final IMessageResolver msgResolver;
private final Function<STrades, IValidator> stradesValidator;
private final KafkaSender kafkaSender;
private static final DateTimeFormatter contractFormatter = DateTimeFormatter.ofPattern("ddMMyy");
@Autowired
public ExecutionCurrencyComponent(ImdgProvider imdgProvider, Producer<String, Object> kafka,
@Qualifier("sTradesValidatorCurrency") Function<STrades, IValidator> stradesValidator,
@Qualifier("kafkaSenderWithoutRequestInfo") KafkaSender kafkaSender,
IMessageResolver msgResolver) {
this.sTradeImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_STrades, STrades.class);
this.executionCurrencyImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionCurrency, ExecutionCurrency.class);
this.listingImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Listing, Listing.class);
this.stradesValidator = stradesValidator;
this.kafkaSender = kafkaSender;
this.msgResolver = msgResolver;
// 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() {
LocalDate today = LocalDate.now();
log.debug("Start check new S_TRADE at {}", today);
//выбираем STrades на сегодня с правильным section
Collection<STrades> sTrades = sTradeImdg.getCollectionObjectsByFieldValues(Map.of(
"tradeDate", today,
"section", Section.CURR.getKey()
));
log.info("Found {} s_trade for today", sTrades.size());
if (sTrades.isEmpty()) {
logError(ClearingError.NewDealsNotFound);
return;
}
//убираем уже добавленные в ExecutionCurrency
sTrades.removeIf(sTrd -> {
Side sTrdSide = IEnumKey.getEnumByKey(Side.class, sTrd.getOperation());
if (sTrdSide == null || !(sTrdSide.equals(Side.BUY) || sTrdSide.equals(Side.SELL))) {
log.warn("cannot parse strade.tradeNum={} strade.operation {}", sTrd.getTradeNum(), sTrd.getOperation());
return true;
}
// MoneyFlowSide excDepSide = sTrdSide == Side.BUY ? MoneyFlowSide.BUY : MoneyFlowSide.SELL;
return executionCurrencyImdg.getFirstObjectByFieldValues(
Map.of("tradingDate", sTrd.getTradeDate(),
"exchangeExecutionId", sTrd.getTradeNum(),
"side", sTrdSide.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;
}
ExecutionCurrency newEC;
try {
newEC = createExecutionCurrency(sTrd, validator);
executionCurrencyImdg.insert(newEC);
sendNotification(newEC);
log.debug("New executionCurrency.id={} was created.", newEC.getId());
} catch (ClearingException ce) {
auditMessage(ce);
} catch (Exception e) {
log.error("When create new ExecutionCurrency by STrade[{}] error: {}", sTrd.getId(), ExceptionUtils.getStackTrace(e));
}
}
Long newMaxTradeNum = sTrades.stream().mapToLong(STrades::getTradeNum).max().orElse(0); // 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(ExecutionCurrency forED) throws ClearingException {
final String destination = Consts.REGISTRY_DEAL_REGISTER_NEW;
DealRegisterNewRequest requestPayload = new DealRegisterNewRequest();
requestPayload.setExecutionId(forED.getId());
requestPayload.setExchangeExecutionId(forED.getExchangeExecutionId());
requestPayload.setExecutionType(ExecutionType.currency);
Long generatedRequestId = kafkaSender.sendRequestToQueue(destination, requestPayload);
if (generatedRequestId == null) {
throw new RuntimeException("failed to send request to " + destination);
}
}
protected ExecutionCurrency createExecutionCurrency(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);
Listing listing = listingImdg.getFirstObjectByFieldValues(Map.of("securityId", security.getId()));
ExecutionCurrency eCurrency = new ExecutionCurrency();
final Instant now = Instant.now();
eCurrency.setCreated(now);
eCurrency.setTradingDate(sTrades.getTradeDate());
eCurrency.setClearingDate(TimeUtil.toLocalDate(now));
eCurrency.setExchangeExecutionId(sTrades.getTradeNum());
eCurrency.setExchangeExecutionTime(sTrades.getTradeDateTime());
eCurrency.setExchangeExecutionMicroseconds(sTrades.getTradeDateTime()); // todo проверить это дата+время или нет
eCurrency.setPartyTradingClearingRegistryId(rgstr.getId()); // setTradingClearingRegistryId
eCurrency.setPartyTradingClearingRegistry(sTrades.getAccount()); // todo уточнить rgstr.code/strades.account?
eCurrency.setMarket(sTrades.getClassCode());
// { todo нет поля Description и не понятно как искать Market
// Market market = market=market.code ;
// eCurrency.setDescription(market.getDescription());
// }
eCurrency.setPrice(sTrades.getPrice());
eCurrency.setLots(sTrades.getQty());
eCurrency.setSettlementAmount(sTrades.getValue());
if (listing != null && listing.getLotSize() != null && eCurrency.getLots() != null) {
BigDecimal lots = eCurrency.getLots();
BigDecimal listingLotSize = listing.getLotSize();
if (listingLotSize != null && lots != null) {
eCurrency.setQuantity(lots.multiply(listingLotSize));
}
}
//todo fill eCurrency.setSessionId();
// String companyContract = null;
{
Side sTradeSide = IEnumKey.getEnumByKey(Side.class, sTrades.getOperation());
if (sTradeSide != null) {
eCurrency.setSide(sTradeSide.getKey());
// switch (sTradeSide) {
// case BUY -> {
// companyContract = sTrades.getFirmId();
// eCurrency.setSide(?MoneyFlowSide.BUY.getKey());
// }
// case SELL -> {
// companyContract = sTrades.getCpFirmId();
// eCurrency.setSide(?MoneyFlowSide.SELL.getKey());
// }
// }
}
}
eCurrency.setCurrencyCode(sTrades.getSettleCurrency());//todo verify TZ
// eCurrency.setSettlementCurrency(sTrades.getSettleCurrency());
eCurrency.setCompanyId(company.getId());
//todo fill eCurrency.setSettlementOrganization();
//todo fill eCurrency.setCoverageStatus(); "атус достаточности обеспечения" будет ли на следующих этапах?
eCurrency.setSettlementCode(sTrades.getSettleCode()); // todo не по заданию, но наверное так
eCurrency.setSettlementDate(sTrades.getSettleDate());
eCurrency.setSecurityName(security.getFullName());
eCurrency.setSecuritySymbol(security.getSecuritySymbol());
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.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());
}
}
return eCurrency;
}
private Optional<Long> stripDaysFromSecCode(String secCode) {
if (secCode == null || secCode.length() <= 7) {
return Optional.empty();
}
String digitsFromSecCode = secCode.substring(7).replaceAll("[^\\d]", "");
if (digitsFromSecCode.length() == 0) {
return Optional.empty();
}
try {
Long days = Long.valueOf(digitsFromSecCode);
return Optional.of(days);
} catch (Throwable e) {
return Optional.empty();
}
}
private void logError(IEnumId subject, Object... args) {
log.info("{}", msgResolver.resolve(new EnumMessage(subject, args)));
}
private void logError(Long sTradeId, EnumMessage msg) {
log.warn("sTrade id={} {}", sTradeId, msgResolver.resolve(msg));
}
private Optional<TradingClearingRegistry> searchTcrByStrades(STrades sTrades) {
IValidator counterValidator = stradesValidator.apply(sTrades);
Optional<EnumMessage> 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));
}
}

View file

@ -3,6 +3,7 @@ package ru.spcex.clearing.service.validation;
import ru.clearing.classes.statics.data.company.Company;
import ru.clearing.classes.statics.data.misc.STrades;
import ru.clearing.classes.statics.data.registry.TradingClearingRegistry;
import ru.clearing.classes.statics.data.security.CurrencyPairSecurity;
import ru.clearing.classes.statics.data.security.MoneyMarketSecurity;
import ru.clearing.classes.statics.data.security.Security;
import ru.spcex.clearing.error.ClearingError;
@ -62,6 +63,35 @@ public enum STradesValidationRule implements IValidationRule<ImdgValidationConte
return empty();
}
},
SecurityPresentCurrencyPair() {
@Override
public Optional<EnumMessage> validate(ImdgValidationContext<STrades> context) {
STrades validatedObject = context.getValidatedObject();
if (TextUtil.isEmpty(validatedObject.getSecCode())) {
return of(ClearingError.SecurityNotFound, validatedObject.getSecCode());
}
String secCode = validatedObject.getSecCode().trim();
SecuritySelector<Security> slctr = new SecuritySelector<>(
context.obtainMap(IMDGDistributedNames.Map_FixedIncomeSecurity, Security.class),
context.obtainMap(IMDGDistributedNames.Map_MoneyMarketSecurity, Security.class),
context.obtainMap(IMDGDistributedNames.Map_EquitySecurity, Security.class),
context.obtainMap(IMDGDistributedNames.Map_CurrencyPairSecurity, Security.class)
);
Security security = slctr.selectSecurityBySymbol(secCode);
/*
Imdg<CurrencyPairSecurity> securityImdg = context.obtainMap(IMDGDistributedNames.Map_CurrencyPairSecurity, CurrencyPairSecurity.class);
String secCode = validatedObject.getSecCode();
Security security = securityImdg.getFirstObjectByFieldValues(Map.of("securitySymbol", secCode));
*/
if (security == null) {
return of(ClearingError.SecurityNotFound, secCode);
}
context.storeObject(ValidationStored.STradesSecurity, security);
return empty();
}
},
// 2.2. Найти в company запись, у которой company.tradingCode=sTrades.firmId.
// Если такой записи нет, записать в лог ошибку (5410) "Компания %s не найдена".
CompanyPresent() {

View file

@ -13,6 +13,7 @@ import org.springframework.test.context.junit.jupiter.SpringExtension;
import ru.spcex.clearing.config.*;
import ru.spcex.clearing.service.builder.LiabilitiesClaimsAssetsCreator;
import ru.spcex.clearing.service.builder.LiabilitiesClaimsMoneyCreator;
import ru.spcex.clearing.service.execution.ExecutionCurrencyComponent;
import ru.spcex.clearing.service.execution.ExecutionDepositComponent;
import ru.spcex.clearing.service.order.ExecutionDepositSorter;
import ru.spcex.clearing.service.order.PaymentInstructionSorter;
@ -29,6 +30,7 @@ import static org.mockito.Mockito.spy;
VerificationResultComponent.class,
ExecutionDepositComponent.class,
ExecutionDepositSorter.class,
ExecutionCurrencyComponent.class,
LiabilitiesClaimsAssetsCreator.class,
LiabilitiesClaimsMoneyCreator.class,
EventsReceiver.class,

View file

@ -0,0 +1,95 @@
package ru.spcex.clearing.service.execution;
import org.apache.kafka.clients.consumer.MockConsumer;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import ru.clearing.classes.statics.data.account.Account;
import ru.clearing.classes.statics.data.company.Company;
import ru.clearing.classes.statics.data.execution.ExecutionCurrency;
import ru.clearing.classes.statics.data.misc.Listing;
import ru.clearing.classes.statics.data.misc.STrades;
import ru.clearing.classes.statics.data.security.Security;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
import ru.spcex.clearing.service.AbstractClearingTest;
import ru.spcex.clearing.service.EventsReceiver;
import ru.spcex.platform.enumeration.Task;
import ru.spcex.platform.imdg.api.Imdg;
import javax.annotation.PostConstruct;
import java.time.Instant;
import java.time.LocalDate;
import java.util.concurrent.atomic.AtomicInteger;
import static ru.spcex.clearing.utils.TestUtils.addRecordToKafka;
import static ru.spcex.clearing.utils.TestUtils.getJsonStringForUPDATE;
class ExecutionCurrencyComponentTest extends AbstractClearingTest {
private static final int PARTITION = 0;
private static final AtomicInteger currentInteger = new AtomicInteger(1);
private static final String TOPIC = Task.getOfTrades.topic();
@Autowired
EventsReceiver eventsReceiver;
@Autowired
ExecutionCurrencyComponent executionCurrencyComponent;
private Instant todayInstant;
private Imdg<STrades> sTradeImdg;
private Imdg<Security> securityImdg;
private Imdg<ExecutionCurrency> executionCurrencyImdg;
private Imdg<Company> companyImdg;
private Imdg<Account> accountImdg;
private Imdg<Listing> listingImdg;
@PostConstruct
protected void init() {
super.init();
// executionCurrencyComponent.resetTradingDay();
this.todayInstant = Instant.now(); //executionCurrencyComponent.tradingDay;
this.sTradeImdg = hazelcastServiceTest.getImdg(IMDGDistributedNames.Map_STrades, STrades.class);
this.securityImdg = hazelcastServiceTest.getImdg(IMDGDistributedNames.Map_Security, Security.class);
this.executionCurrencyImdg = hazelcastServiceTest.getImdg(IMDGDistributedNames.Map_ExecutionCurrency, ExecutionCurrency.class);
this.companyImdg = hazelcastServiceTest.getImdg(IMDGDistributedNames.Map_Company, Company.class);
this.accountImdg = hazelcastServiceTest.getImdg(IMDGDistributedNames.Map_Account, Account.class);
this.listingImdg = hazelcastServiceTest.getImdg(IMDGDistributedNames.Map_Listing, Listing.class);
}
/**
* {@link ExecutionCurrencyComponent#processNewTS()} ()} <br>
* Тест проверяет создание сущностей
* {@link ExecutionCurrency}
* в Hazelcast при передаче из Apache Kafka.<br>
*/
@Test
void processNewTS() {
String secCode = "SecCode";
String operation = "oper";
STrades sTrades = new STrades();
Long exchangeExecutionId = 1221L;
LocalDate today = LocalDate.now();
sTrades.setTradeDateTime(todayInstant);
sTrades.setSecCode(secCode);
sTrades.setOperation(operation);
sTrades.setSection("CURR");
sTradeImdg.insert(sTrades);
Security security = new Security();
security.setSecuritySymbol(secCode);
securityImdg.insert(security);
ExecutionCurrency executionCurrency = new ExecutionCurrency();
executionCurrency.setSide(operation);
executionCurrency.setClearingDate(today);
executionCurrency.setExchangeExecutionId(exchangeExecutionId);
executionCurrencyImdg.insert(executionCurrency);
int times = currentInteger.getAndIncrement();
addRecordToKafka((MockConsumer) eventsReceiver.getConsumer(), TOPIC,
PARTITION, times, getJsonStringForUPDATE(new LauncherCommandRequest(), 1L));
//waiting for kafka producer send message (finale event)
// verify(mockProducer, timeout(30_000L).times(times))
// .send(producerRecord.capture());
}
}

View file

@ -4,7 +4,8 @@ import ru.spcex.platform.utils.enumeration.IEnumKey;
public enum Section implements IEnumKey {
MKR("MKR"),
FOND("FOND");
FOND("FOND"),
CURR("CURR");
Section(String key) {
this.key = key;

View file

@ -1,5 +1,5 @@
package ru.spcex.clearing.platform.messaging.domain.cud.registry;
public enum ExecutionType {
deposit, fond;
deposit, fond, currency;
}