From 81dbf7bf2c9a5ac2536d4deb17f7f7e9339dc68c Mon Sep 17 00:00:00 2001 From: etreshenkov Date: Thu, 25 May 2023 13:13:20 +0300 Subject: [PATCH] http://jira.mfd.msk:8088/browse/CLS-290 --- .../statics/data/registry/Registry.java | 6 +- .../clearing/service/EventsReceiver.java | 2 +- .../clearing/service/StatementService.java | 8 +- .../builder/PaymentInstructionBuilder.java | 74 +++++++++++++++++++ .../service/executors/AbstractExecutor.java | 4 +- .../service/executors/Sdf01Executor.java | 13 +++- .../service/executors/Sdf57Executor.java | 24 +++++- .../session/stage/impl/BalanceRevise.java | 43 +++++------ .../platform/messaging/domain/Consts.java | 2 +- .../clearing/ContinueSessionBnRequest.java | 16 ++++ 10 files changed, 151 insertions(+), 41 deletions(-) create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/builder/PaymentInstructionBuilder.java create mode 100644 platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/clearing/ContinueSessionBnRequest.java diff --git a/clearing-parent/classes/src/main/java/ru/clearing/classes/statics/data/registry/Registry.java b/clearing-parent/classes/src/main/java/ru/clearing/classes/statics/data/registry/Registry.java index 93836225e..dd68a819e 100644 --- a/clearing-parent/classes/src/main/java/ru/clearing/classes/statics/data/registry/Registry.java +++ b/clearing-parent/classes/src/main/java/ru/clearing/classes/statics/data/registry/Registry.java @@ -2,10 +2,8 @@ package ru.clearing.classes.statics.data.registry; import ru.clearing.classes.ConstSerializable; import ru.clearing.classes.objects.BusinessObject; -import ru.spcex.platform.classes.base.SpcexObjectBase; import java.math.BigDecimal; -import java.time.Instant; import java.time.LocalDate; /** @@ -414,9 +412,9 @@ public class Registry extends BusinessObject implements Cloneable { } @Override - public Object clone() { + public Registry clone() { try { - return super.clone(); + return (Registry) super.clone(); } catch (CloneNotSupportedException e) { throw new RuntimeException("never", e); // it is Cloneable } 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 3e51ed48e..dbe362fe4 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 @@ -63,7 +63,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean { .forDestination(Task.startOfClearing.topic(), callbacks::put); callback(Object.class) .setConsumer(primaryAuctionBnSession::continueSession) - .forDestination(Consts.SDF57_PROCESS, callbacks::put); + .forDestination(Consts.CONTINUE_SESSION_BN, callbacks::put); callback(Object.class) .setConsumer(event -> clearingService.executeSTrade()) diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java index dd42e9530..208158ed2 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java @@ -14,7 +14,6 @@ import ru.spcex.clearing.platform.messaging.domain.BaseRequest; import ru.spcex.clearing.platform.messaging.domain.Consts; import ru.spcex.clearing.platform.messaging.domain.cud.account.sdf01.AccountSdf01Request; import ru.spcex.clearing.platform.messaging.domain.cud.account.sdf01.AccountSdfRequestPart; -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.service.QueueConsumer; import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; @@ -81,11 +80,8 @@ public class StatementService extends QueueConsumer implements InitializingBean Result res = service.execute(sdfGroup, statementRequest); if (res.getAccountRequests().size() != 0) { kafkaSender.sendRequestToQueue(Consts.ACCOUNT_NEW_SDF01, createAccountsRequest(statementRequest.getGroupId(), res.getAccountRequests())); - } else if (service.isNeedToSendCommandToExport()) { - ExportToFileRequest exportRequest = new ExportToFileRequest(); - exportRequest.setSdfGroupId(res.getGenerationId()); - exportRequest.setNameOfTable(service.exportTableName()); - kafkaSender.sendRequestToQueue(Consts.EXPORT_PROCESS, exportRequest); + } else if (service.isNeedToSendCommand()) { + service.sendCommand(kafkaSender, res); } } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/builder/PaymentInstructionBuilder.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/builder/PaymentInstructionBuilder.java new file mode 100644 index 000000000..ff1f9ecde --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/builder/PaymentInstructionBuilder.java @@ -0,0 +1,74 @@ +package ru.spcex.clearing.service.builder; + +import ru.clearing.classes.statics.data.account.Account; +import ru.clearing.classes.statics.data.company.Company; +import ru.clearing.classes.statics.data.company.CompanySymbols; +import ru.clearing.classes.statics.data.payment.PaymentInstruction; +import ru.clearing.classes.statics.data.registry.Registry; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.platform.enumeration.CompanySymbol; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; + +import java.util.Map; + +public class PaymentInstructionBuilder { + + private Registry registry; + private ImdgProvider imdgProvider; + + public static PaymentInstructionBuilder builder(ImdgProvider imdgProvider, Registry registry) { + return new PaymentInstructionBuilder(imdgProvider, registry); + } + + private PaymentInstructionBuilder(ImdgProvider imdgProvider, Registry registry) { + this.imdgProvider = imdgProvider; + this.registry = registry; + } + + public PaymentInstruction build() { + PaymentInstruction paymentInstruction = new PaymentInstruction(); + paymentInstruction.setSenderId(registry.getCompanyId()); + paymentInstruction.setAddresseeId(registry.getCounterPartyId()); + + Imdg companySymbolsImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_CompanySymbols, CompanySymbols.class); + CompanySymbols companySymbols = companySymbolsImdg.getSingleObjectByFieldValues(Map.of( + "companyId", registry.getCompanyId(), + "companySymbol", CompanySymbol.BIC.getKey()) + ); + CompanySymbols counterCompanySymbols = companySymbolsImdg.getSingleObjectByFieldValues(Map.of( + "companyId", registry.getCounterPartyId(), + "companySymbol", CompanySymbol.BIC.getKey()) + ); + + if (companySymbols != null) { + paymentInstruction.setPayeeBic(companySymbols.getCompanySymbolValue()); + } + if (counterCompanySymbols != null) { + paymentInstruction.setAdresseeBic(counterCompanySymbols.getCompanySymbolValue()); + } + + Imdg companyImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Company, Company.class); + Company company = companyImdg.getSingleObjectByFieldValues(Map.of("companyId", registry.getCompanyId())); + if (company != null) { + paymentInstruction.setPayeeBankName(company.getShortName()); + paymentInstruction.setAddresseeBankName(company.getShortName()); + } + paymentInstruction.setSettlementDate(registry.getSettlementDate()); + paymentInstruction.setCreditLeg_amount(registry.getBalance()); + paymentInstruction.setDebitLeg_amount(registry.getBalance()); + paymentInstruction.setCreditLeg_accountId(registry.getAccountId()); + + Imdg accountImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Account, Account.class); + Account account = accountImdg.getSingleObjectByID(registry.getAccountId()); + if (account != null) { + paymentInstruction.setCreditLeg_account(account.getAccount()); + } + + paymentInstruction.setDebitLeg_accountId(registry.getAccountId()); + +// Account account = accountImdg.getSingleObjectByID(registry.get()); +// paymentInstruction.setDebitLeg_account(); + return paymentInstruction; + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/AbstractExecutor.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/AbstractExecutor.java index b1bc4cddf..afead4914 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/AbstractExecutor.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/AbstractExecutor.java @@ -1,6 +1,7 @@ package ru.spcex.clearing.service.executors; import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.clearing.service.model.Result; import java.util.Collection; @@ -8,5 +9,6 @@ import java.util.Collection; public abstract class AbstractExecutor { public abstract Result execute(Collection sdf, StatementRequest statementRequest); public abstract String exportTableName(); - public abstract boolean isNeedToSendCommandToExport(); + public abstract boolean isNeedToSendCommand(); + public abstract void sendCommand(KafkaSender kafkaSender, Result result); } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf01Executor.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf01Executor.java index e2b30c58c..8f18d90cc 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf01Executor.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf01Executor.java @@ -11,8 +11,11 @@ import ru.clearing.classes.statics.data.sdf.SDf02; import ru.clearing.classes.statics.data.statement.Statement; import ru.spcex.clearing.error.ClearingError; import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.platform.messaging.domain.Consts; import ru.spcex.clearing.platform.messaging.domain.cud.account.sdf01.AccountSdfRequestPart; +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.service.sender.KafkaSender; import ru.spcex.clearing.service.LoggingService; import ru.spcex.clearing.service.model.Result; import ru.spcex.clearing.service.validation.ValidationStored; @@ -65,10 +68,18 @@ public class Sdf01Executor extends AbstractExecutor { } @Override - public boolean isNeedToSendCommandToExport() { + public boolean isNeedToSendCommand() { return true; } + @Override + public void sendCommand(KafkaSender kafkaSender, Result result) { + ExportToFileRequest exportRequest = new ExportToFileRequest(); + exportRequest.setSdfGroupId(result.getGenerationId()); + exportRequest.setNameOfTable(exportTableName()); + kafkaSender.sendRequestToQueue(Consts.EXPORT_PROCESS, exportRequest); + } + public Result execute(Collection sdf, StatementRequest statementRequest) { Result result = new Result(); Long generationIdForGroup = imdgProvider.getImdgIdGenerator().nextId(); diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf57Executor.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf57Executor.java index 21fd8fa77..7be4a4e1b 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf57Executor.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf57Executor.java @@ -14,7 +14,10 @@ import ru.clearing.classes.statics.data.security.Security; import ru.clearing.classes.statics.data.statement.Statement; import ru.spcex.clearing.error.ClearingError; import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.platform.messaging.domain.Consts; import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.clearing.ContinueSessionBnRequest; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.clearing.service.LoggingService; import ru.spcex.clearing.service.model.Result; import ru.spcex.clearing.service.validation.ValidationStored; @@ -78,9 +81,16 @@ public class Sdf57Executor extends AbstractExecutor { } @Override - public boolean isNeedToSendCommandToExport() { - return false; - }//no need... + public boolean isNeedToSendCommand() { + return true; + } + + @Override + public void sendCommand(KafkaSender kafkaSender, Result result) { + ContinueSessionBnRequest continueSessionBn = new ContinueSessionBnRequest(); + continueSessionBn.setGenerationId(result.getGenerationId()); + kafkaSender.sendRequestToQueue(Consts.CONTINUE_SESSION_BN, continueSessionBn); + } //V - Изменение statement по sDf57 // @@ -145,7 +155,13 @@ public class Sdf57Executor extends AbstractExecutor { }; Consumer create = (dsgn) -> { Registry registry = createRegistryByStatement(stmt, company, account, dsgn); + Registry registryF = registry.clone(); + registryF.setRegistryUnit(RegistryUnit.F.getKey()); + Registry registryB = registry.clone(); + registryB.setRegistryUnit(RegistryUnit.B.getKey()); registryImdg.insert(registry); + registryImdg.insert(registryF); + registryImdg.insert(registryB); }; //todo понять что происходит со инициатором/контрагентом, особенно если у нас только один Statement @@ -181,7 +197,7 @@ public class Sdf57Executor extends AbstractExecutor { Statement statement = new Statement(); statement.setAddresseeId(companyDeb.getId()); statement.setSenderId(Sender.Prc.getId()); - statement.setStatementType(StatementType.full.getKey()); + statement.setStatementType(StatementType.incr.getKey()); statement.setContract(getContractFromSpecif(sdf57.getSpecif())); statement.setAccountId(accountDeb.getId()); statement.setAccount(accountDeb.getAccount()); diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/BalanceRevise.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/BalanceRevise.java index d962c1355..94346090e 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/BalanceRevise.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/BalanceRevise.java @@ -19,6 +19,7 @@ import ru.spcex.platform.enumeration.*; 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.specific.RegistryCodeSqlBuilder; import ru.spcex.platform.imdg.api.predicate.specific.StatementRevisePredicate; import ru.spcex.platform.utils.enumeration.IEnumKey; @@ -93,32 +94,31 @@ public class BalanceRevise implements ISessionStage { Collection stmts = statementImdg.getCollectionObjectsBySQL(statementSQL); for (Statement stmt : stmts) { - String unformatted = "registry_designation = '%s' " + - "and registry_instrument_type = '%s' " + - "and registry_unit = '%s' " + - "and account = '%s' " + - "and securityId = %d"; + String unformatted = "%s and account = '%s' and securityId = %d"; + + RegistryTradingParams registryAMT = new RegistryTradingParams(RegistryDesignation.A, RegistryInstrumentType.M, + null, RegistryUnit.T); + RegistryTradingParams registryAMF = new RegistryTradingParams(RegistryDesignation.A, RegistryInstrumentType.M, + null, RegistryUnit.F); + RegistryTradingParams registryAMB = new RegistryTradingParams(RegistryDesignation.A, RegistryInstrumentType.M, + null, RegistryUnit.B); + String registrySqlAMT = String.format(unformatted, - RegistryDesignation.A.getKey(), - RegistryInstrumentType.M.getKey(), - RegistryUnit.T.getKey(), + RegistryCodeSqlBuilder.getInstance(registryAMT).build(), stmt.getAccount(), stmt.getSecurityId() ); String registrySqlAMF = String.format(unformatted, - RegistryDesignation.A.getKey(), - RegistryInstrumentType.M.getKey(), - RegistryUnit.F.getKey(), + RegistryCodeSqlBuilder.getInstance(registryAMF).build(), stmt.getAccount(), stmt.getSecurityId() ); String registrySqlAMB = String.format(unformatted, - RegistryDesignation.A.getKey(), - RegistryInstrumentType.M.getKey(), - RegistryUnit.B.getKey(), + RegistryCodeSqlBuilder.getInstance(registryAMB).build(), stmt.getAccount(), stmt.getSecurityId() ); + Registry rgsAMT = registryImdg.getSingleObjectBySQL(registrySqlAMT); Registry rgsAMF = registryImdg.getSingleObjectBySQL(registrySqlAMF); Registry rgsAMB = registryImdg.getSingleObjectBySQL(registrySqlAMB); @@ -152,14 +152,11 @@ public class BalanceRevise implements ISessionStage { InOutSDfType.type1.getKey(), OperationStatus.Pending.getKey(), StatementType.full.getKey()); Collection stmts = statementImdg.getCollectionObjectsBySQL(statementSQL); for (Statement stmt : stmts) { - String registrySqlAMT = String.format("registry_designation = '%s' " + - "and registry_instrument_type = '%s' " + - "and registry_unit = '%s' " + - "and account = '%s' " + - "and securityId = %d", - RegistryDesignation.A.getKey(), - RegistryInstrumentType.M.getKey(), - RegistryUnit.T.getKey(), + + RegistryTradingParams registryAMT = new RegistryTradingParams(RegistryDesignation.A, RegistryInstrumentType.M, + null, RegistryUnit.T); + String registrySqlAMT = String.format("%s and account = '%s' and securityId = %d", + RegistryCodeSqlBuilder.getInstance(registryAMT).build(), stmt.getAccount(), stmt.getSecurityId() ); @@ -175,7 +172,7 @@ public class BalanceRevise implements ISessionStage { } private void newSDf56(Statement statement) { - log.debug("GALB request received; creating sdf56"); + log.debug("creating sdf56"); SDf56 sDf56 = new SDf56(); sDf56.setNumber(idGenerator.nextId().toString()); Instant now = Instant.now(); 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 e0197d92c..a897be3b6 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 @@ -121,7 +121,7 @@ public interface Consts { String SDF53_PROCESS = "sdf53-process"; String SDF54_PROCESS = "sdf54-process"; String SDF56_PROCESS = "sdf56-process"; - String SDF57_PROCESS = "sdf57-process"; + String CONTINUE_SESSION_BN = "sdf57-process"; String REVISE_PROCESS = "revise-process"; String EXPORT_PROCESS = "export-process"; String EXPORT_COMPLETED = "export_completed"; diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/clearing/ContinueSessionBnRequest.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/clearing/ContinueSessionBnRequest.java new file mode 100644 index 000000000..e02afec3c --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/clearing/ContinueSessionBnRequest.java @@ -0,0 +1,16 @@ +package ru.spcex.clearing.platform.messaging.domain.cud.clearing; + +import com.fasterxml.jackson.annotation.JsonProperty; + +public class ContinueSessionBnRequest { + @JsonProperty + private Long generationId; + + public Long getGenerationId() { + return generationId; + } + + public void setGenerationId(Long generationId) { + this.generationId = generationId; + } +}