From 5a772d40a3db610084c10c283039a0923e479177 Mon Sep 17 00:00:00 2001 From: ialbert Date: Thu, 2 Feb 2023 16:14:06 +0300 Subject: [PATCH] http://jira.mfd.msk:8088/browse/CLS-189 --- .../balance/service/StatementService.java | 4 ++ .../ru/spcex/clearing/service/Clearing.java | 48 +++++++++++-- .../clearing/service/ClearingService.java | 10 +++ .../clearing/service/EventsReceiver.java | 4 ++ .../LiabilitiesClaimsAssetsCreator.java | 69 +++++++++++++++++++ .../platform/messaging/domain/Consts.java | 1 + .../domain/cud/common/CommonIdRequest.java | 17 +++++ 7 files changed, 149 insertions(+), 4 deletions(-) create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/builder/LiabilitiesClaimsAssetsCreator.java create mode 100644 platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/common/CommonIdRequest.java diff --git a/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/service/StatementService.java b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/service/StatementService.java index 9bedcb348..71d790c18 100644 --- a/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/service/StatementService.java +++ b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/service/StatementService.java @@ -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 updateAccBalanceReq) { accountBalanceService.updateAccountBalanceByClearing(updateAccBalanceReq.getRequestPayload()); + CommonIdRequest commonIdRequest = new CommonIdRequest(); + commonIdRequest.setId(updateAccBalanceReq.getId()); + kafkaReqProducer.sendRequestToQueue(Consts.CONTINUE_CLEARING, commonIdRequest); } private void process(BaseRequest systemRequest) { diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/Clearing.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/Clearing.java index 935866921..4e2136cc0 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/Clearing.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/Clearing.java @@ -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 executionDepositImdg; + private final Imdg liabilitiesClaimsAssetsImdg; private final Imdg clearingCategoryImdg; private final ExecutionDepositSorter sorter; private final BiFunction 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 validation) { + BiFunction 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; + } } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/ClearingService.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/ClearingService.java index 18079d147..0dc6ac3b9 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/ClearingService.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/ClearingService.java @@ -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 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(() -> { diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/EventsReceiver.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/EventsReceiver.java index 66b6c8f47..9eb639b9c 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/EventsReceiver.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/EventsReceiver.java @@ -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); diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/builder/LiabilitiesClaimsAssetsCreator.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/builder/LiabilitiesClaimsAssetsCreator.java new file mode 100644 index 000000000..cf4bcd74b --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/builder/LiabilitiesClaimsAssetsCreator.java @@ -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 accountImdg; + private final Imdg 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; + } +} diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java index fd95ea1ad..97ed7998a 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java @@ -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"; diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/common/CommonIdRequest.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/common/CommonIdRequest.java new file mode 100644 index 000000000..57c9ea6fd --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/common/CommonIdRequest.java @@ -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; + } +}