ialbert 2023-02-02 16:14:06 +03:00
parent c59d010444
commit 5a772d40a3
7 changed files with 149 additions and 4 deletions

View file

@ -18,6 +18,7 @@ import ru.spcex.clearing.platform.messaging.domain.cud.account.sdf01.AccountSdfR
import ru.spcex.clearing.platform.messaging.domain.cud.balance.AccountBalanceClearingRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.balance.ExportToFileRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonIdRequest;
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.platform.classes.base.interfaces.WithAccount;
@ -70,6 +71,9 @@ public class StatementService extends QueueConsumer implements InitializingBean
private void accountBalanceClearingUpdate(BaseRequest<AccountBalanceClearingRequest> updateAccBalanceReq) {
accountBalanceService.updateAccountBalanceByClearing(updateAccBalanceReq.getRequestPayload());
CommonIdRequest commonIdRequest = new CommonIdRequest();
commonIdRequest.setId(updateAccBalanceReq.getId());
kafkaReqProducer.sendRequestToQueue(Consts.CONTINUE_CLEARING, commonIdRequest);
}
private void process(BaseRequest<StatementRequest> systemRequest) {

View file

@ -7,17 +7,21 @@ import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.execution.ExecutionDeposit;
import ru.clearing.classes.statics.data.generated.ClearingMemberCategory;
import ru.clearing.classes.statics.data.liabilities.LiabilitiesClaimsAssets;
import ru.clearing.classes.statics.data.payment.PaymentInstruction;
import ru.spcex.clearing.error.ClearingErrorInternal;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.service.order.ExecutionDepositSorter;
import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.balance.AccountBalanceClearingRequest;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.clearing.service.builder.LiabilitiesClaimsAssetsCreator;
import ru.spcex.clearing.service.builder.PaymentInstructionCreator;
import ru.spcex.clearing.service.order.ExecutionDepositSorter;
import ru.spcex.platform.enumeration.Allowed;
import ru.spcex.platform.enumeration.ClearingCategory;
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.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.log.ExceptionUtils;
import ru.spcex.platform.utils.validation.IValidator;
import java.time.Instant;
@ -28,11 +32,15 @@ import java.util.function.BiFunction;
public class Clearing {
private final Logger log = LoggerFactory.getLogger(getClass());
private final Imdg<ExecutionDeposit> executionDepositImdg;
private final Imdg<LiabilitiesClaimsAssets> liabilitiesClaimsAssetsImdg;
private final Imdg<ClearingMemberCategory> clearingCategoryImdg;
private final ExecutionDepositSorter sorter;
private final BiFunction<ClearingCategory, ExecutionDeposit, IValidator> validation;
private final PaymentInstructionCreator paymentInstructionCreator;
private final ImdgId idProvider;
private final KafkaSender kafka;
private final LiabilitiesClaimsAssetsCreator lbltsClmsAssetsCreator;
private Long lastAccBalanceResponseId = 0L;
//при прохождении по выгруженным ExecutionDeposit, ошибочные статусы проставляются для
//контр сделок. В таком случае, в коллекции хранятся не синхронизированные с IMDG ExecutionDeposit
//для которых статус должен быть DENIED
@ -43,13 +51,16 @@ public class Clearing {
@Autowired
public Clearing(ImdgProvider imdgProvider,
ExecutionDepositSorter sorter, PaymentInstructionCreator paymentInstructionCreator,
BiFunction<ClearingCategory, ExecutionDeposit, IValidator> validation) {
BiFunction<ClearingCategory, ExecutionDeposit, IValidator> validation, KafkaSender kafka, LiabilitiesClaimsAssetsCreator liabilitiesClaimsAssetsCreator) {
this.executionDepositImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionDeposit, ExecutionDeposit.class);
this.clearingCategoryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ClearingMemberCategory, ClearingMemberCategory.class);
this.liabilitiesClaimsAssetsImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_LiabilitiesClaimsAssets, LiabilitiesClaimsAssets.class);
this.paymentInstructionCreator = paymentInstructionCreator;
this.idProvider = imdgProvider.getImdgIdGenerator();
this.sorter = sorter;
this.validation = validation;
this.kafka = kafka;
this.lbltsClmsAssetsCreator = liabilitiesClaimsAssetsCreator;
}
public void startClearing() {
@ -134,7 +145,36 @@ public class Clearing {
if (!Allowed.ALLOWED.getKey().equals(execDeposit.getCoverageStatus())) {
continue;
}
//send request to kafka, wait for a reply
synchronized (this) {
AccountBalanceClearingRequest request = new AccountBalanceClearingRequest();
request.setAccountId(execDeposit.getAccountId());
request.setCompanyId(execDeposit.getCompanyId());
request.setFirstLegAmount(execDeposit.getFirstLegAmount());
Long sentRequestId = kafka.sendRequestToQueue(Consts.BALANCE_ACCOUNT_UPDATE, request);
try {
wait(5000);
//todo check that request Id was processed succesfully
//если на этом месте произошла ошибка - непонятно как обрабатывать
//т.к. этом может быть уже вторая часть сделки, для первой accountBalance уже был обновлен
//либо вне зависимости от того какая часть сделки, accountBalance на самом деле мог быть
//обновлен, и была проблема была с сетью/недоступностью кафки etc.
} catch (InterruptedException e) {
log.error(ExceptionUtils.getStackTrace(e));
}
if (!Objects.equals(lastAccBalanceResponseId, sentRequestId)) {
log.error("FATAL: last response from AccountBalance update from Kafka ID was {}, but sent ID was {}",
lastAccBalanceResponseId,
sentRequestId);
continue;
}
}
LiabilitiesClaimsAssets liabilitiesClaimsAssets = lbltsClmsAssetsCreator.createLiabilitiesClaimsAssets(category, execDeposit);
liabilitiesClaimsAssetsImdg.insert(liabilitiesClaimsAssets);
}
}
public void setLastAccBalanceResponseId(Long lastAccBalanceResponseId) {
this.lastAccBalanceResponseId = lastAccBalanceResponseId;
}
}

View file

@ -6,6 +6,8 @@ import org.springframework.beans.factory.DisposableBean;
import org.springframework.beans.factory.annotation.Autowired;
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.platform.utils.log.ExceptionUtils;
import java.util.concurrent.ExecutorService;
@ -85,6 +87,14 @@ public class ClearingService implements DisposableBean {
});
}
public void continueClearing(BaseRequest<CommonIdRequest> event) {
synchronized (clearing) {
CommonIdRequest requestPayload = event.getRequestPayload();
clearing.setLastAccBalanceResponseId(requestPayload.getId());
clearing.notifyAll();
}
}
public void executeSTrade() {
log.info("start STrade check for ExecutionDeposit");
executor.execute(() -> {

View file

@ -5,6 +5,7 @@ import org.springframework.beans.factory.InitializingBean;
import org.springframework.stereotype.Service;
import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.Sdf04Request;
import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonIdRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
import ru.spcex.platform.enumeration.Task;
@ -34,6 +35,9 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
callback(Object.class) //todo check Object suitable
.setConsumer(event -> clearingService.executeClearing())
.forDestination(Task.startOfClearing.topic(), callbacks::put);
callback(CommonIdRequest.class)
.setConsumer(clearingService::continueClearing)
.forDestination(Consts.CONTINUE_CLEARING, callbacks::put);
callback(Object.class)
.setConsumer(event -> clearingService.executeSTrade())
.forDestination(Task.getOfTrades.topic(), callbacks::put);

View file

@ -0,0 +1,69 @@
package ru.spcex.clearing.service.builder;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import ru.clearing.classes.statics.data.account.Account;
import ru.clearing.classes.statics.data.company.Company;
import ru.clearing.classes.statics.data.execution.ExecutionDeposit;
import ru.clearing.classes.statics.data.liabilities.LiabilitiesClaimsAssets;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.platform.enumeration.AccountType;
import ru.spcex.platform.enumeration.ClearingCategory;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.utils.enumeration.IEnumKey;
import java.time.Instant;
import java.time.LocalDate;
@Component
public class LiabilitiesClaimsAssetsCreator {
private final Imdg<Account> accountImdg;
private final Imdg<Company> companyImdg;
@Autowired
public LiabilitiesClaimsAssetsCreator(ImdgProvider imdgProvider) {
this.accountImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Account, Account.class);
this.companyImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Company, Company.class);
}
public LiabilitiesClaimsAssets createLiabilitiesClaimsAssets(ClearingCategory category, ExecutionDeposit executionDeposit) {
LiabilitiesClaimsAssets liabilitiesClaimsAssets = new LiabilitiesClaimsAssets();
liabilitiesClaimsAssets.setCompanyId(executionDeposit.getCompanyId());
liabilitiesClaimsAssets.setAccountId(executionDeposit.getAccountId());
Account account = accountImdg.getSingleObjectByID(executionDeposit.getAccountId());
liabilitiesClaimsAssets.setAccountType(account.getAccountType());
liabilitiesClaimsAssets.setAccount(account.getAccount());
AccountType accType = IEnumKey.getEnumByKey(AccountType.class, account.getAccountType());
boolean iClrn = ClearingCategory.I.equals(category) && AccountType.Clrn.equals(accType);
boolean vInfo = ClearingCategory.V.equals(category) && AccountType.Info.equals(accType);
boolean bClrn = ClearingCategory.B.equals(category) && AccountType.Clrn.equals(accType);
if (iClrn || vInfo) {
liabilitiesClaimsAssets.setLiabilitiesQuantity(executionDeposit.getFirstLegAmount());
}
if (bClrn) {
liabilitiesClaimsAssets.setClaimsQuantity(executionDeposit.getSecondLegAmount());
}
liabilitiesClaimsAssets.setCurrency(executionDeposit.getSettlementCurrency());
liabilitiesClaimsAssets.setSettlementDate(executionDeposit.getFirstLegSettlementDate());
liabilitiesClaimsAssets.setTradingDate(executionDeposit.getTradingDate());
liabilitiesClaimsAssets.setRefundDate(executionDeposit.getSecondLegSettlementDate());
liabilitiesClaimsAssets.setPrice(executionDeposit.getPrice());
liabilitiesClaimsAssets.setSecurityId(executionDeposit.getSecurityId());
Company company = companyImdg.getSingleObjectByID(executionDeposit.getCompanyId());
liabilitiesClaimsAssets.setTradingCode(company.getTradingCode());
liabilitiesClaimsAssets.setClearingCode(company.getClearingCode());
liabilitiesClaimsAssets.setShortName(company.getShortName());
liabilitiesClaimsAssets.setContract(executionDeposit.getSecuritySymbol());
liabilitiesClaimsAssets.setFullName(company.getFullName());
//todo liabilitiesClaimsAssets.setLiabilitiesClaimsMoneyId(company.getFullName());
//fixme liabilitiesClaimsAssets.setClearingStatus();
//todo liabilitiesClaimsAssets.setPaymentId();
//todo liabilitiesClaimsAssets.setRefundPaymentId();
liabilitiesClaimsAssets.setCreated(Instant.now());
liabilitiesClaimsAssets.setClearingDate(LocalDate.now());
return liabilitiesClaimsAssets;
}
}

View file

@ -51,6 +51,7 @@ public interface Consts {
String ACCOUNT_NEW = "account-new";
String BALANCE_ACCOUNT_NEW = "balance-account-new";
String BALANCE_ACCOUNT_UPDATE = "balance-account-update";
String CONTINUE_CLEARING = "continue-clearing";
String LAUNCHER_NEW = "launcher-new";
String NOTIFICATION_NEW = "notification-new";

View file

@ -0,0 +1,17 @@
package ru.spcex.clearing.platform.messaging.domain.cud.common;
import com.fasterxml.jackson.annotation.JsonProperty;
import ru.spcex.platform.classes.base.interfaces.WithId;
public class CommonIdRequest implements WithId {
@JsonProperty
public Long id;
public Long getId() {
return id;
}
public void setId(Long id) {
this.id = id;
}
}