clearing-service http://jira.mfd.msk:8088/browse/CLS-188 STrade Часть I - Изменение executionDeposit при получении новых сделок из ТС
This commit is contained in:
parent
65086fa8af
commit
11d38b3112
4 changed files with 389 additions and 1 deletions
|
|
@ -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;
|
||||
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
}
|
||||
}
|
||||
|
|
@ -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<STrade> sTradeImdg;
|
||||
private Imdg<Security> securityImdg;
|
||||
private Imdg<ExecutionDeposit> executionDepositImdg;
|
||||
private Imdg<Company> companyImdg;
|
||||
private Imdg<Account> accountImdg;
|
||||
private Imdg<Listing> listingImdg;
|
||||
|
||||
private ImdgId idGenerator;
|
||||
Producer<String, Object> kafka;
|
||||
|
||||
Long tradeNum;
|
||||
Instant tradingDay;
|
||||
|
||||
|
||||
@Autowired
|
||||
public ExecutionDepositComponent(ImdgProvider imdgProvider, Producer<String, Object> 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<STrade> 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<String> 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<Security> foundSecurities = securityImdg.getCollectionObjectsByPredicate(allIn);
|
||||
Set<String> 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<String> 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<ExecutionDeposit> 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<String> newSymbolRequest) {
|
||||
final String destination = Consts.DESTINATION_MONEY_MARKET_SECURITY_NEW;
|
||||
List<String> symbolRequests = new ArrayList<>(newSymbolRequest); // чтобы в случае ошибки отобразить номер в логе
|
||||
List<Future<RecordMetadata>> 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<Object> 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<RecordMetadata> send = kafka.send(new ProducerRecord<>(destination, request));
|
||||
sendAll.add(send);
|
||||
}
|
||||
}
|
||||
int i = 0;
|
||||
for (Future<RecordMetadata> 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<Object> request = new BaseRequest<>();
|
||||
request.setId(idGenerator.nextId());
|
||||
request.setActionType(ActionType.NEW);
|
||||
request.setRequestPayload(requestPayload);
|
||||
log.trace("Send to {} new ExecutionDeposit[{}]", destination, forED.getId());
|
||||
Future<RecordMetadata> 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;
|
||||
}
|
||||
|
||||
}
|
||||
|
|
@ -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 * * * * *
|
||||
clearing-service.scheduler.check-payment-instruction=*/5 * * * * *
|
||||
clearing-service.scheduler.check-s-trade=0 1 0 1 * ?
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue