From 0ad02d573b71387b3fbc7da1c7d4f8508e96ed92 Mon Sep 17 00:00:00 2001 From: ialbert Date: Wed, 23 Aug 2023 11:01:30 +0300 Subject: [PATCH] SDF06, SDF57 DM*I --- .../clearing/config/ValidationConfig.java | 1 + .../spcex/clearing/service/AnltSearcher.java | 16 ++++-- .../service/executors/Sdf06Executor.java | 18 +++++-- .../service/executors/Sdf57Executor.java | 49 ++++++++++++++++++- .../integration/GatewayRequestCreator.java | 18 +++++++ .../clearing/service/registry/DmiService.java | 28 +++++++---- .../validation/Sdf06NewValidationRule.java | 10 ++++ 7 files changed, 121 insertions(+), 19 deletions(-) create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/integration/GatewayRequestCreator.java diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/ValidationConfig.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/ValidationConfig.java index 7edd61603..4c6c5a857 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/ValidationConfig.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/ValidationConfig.java @@ -265,6 +265,7 @@ public class ValidationConfig { context.addImdg(IMDGDistributedNames.Map_TradingClearingRegistry, imdgTradingClearingRegistry); context.setLogPrefix(LogPrefixId.INSTANCE); return new ValidatorImpl<>(context, + Sdf06NewValidationRule.Fields, Sdf06NewValidationRule.CompanySymbolAndCompanyPresent, Sdf06NewValidationRule.AccountPresent, Sdf06NewValidationRule.Session, diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/AnltSearcher.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/AnltSearcher.java index be9322c76..40e549235 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/AnltSearcher.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/AnltSearcher.java @@ -77,7 +77,7 @@ public class AnltSearcher { if (company == null) { return new AnltSearch(ClearingError.CompanyNotFoundB, tcr.getCompanyId()); } - return new AnltSearch(infoAcc, company); + return new AnltSearch(infoAcc, company, tcr); } private static String getTkrCodeFromComment(String comment) { @@ -97,18 +97,20 @@ public class AnltSearcher { private EnumMessage error; private Account account; private Company company; + private TradingClearingRegistry tcr; public AnltSearch(IEnumId errSubject, Object... args) { this.error = new EnumMessage(errSubject, args); } - public AnltSearch(Account account, Company company) { + public AnltSearch(Account account, Company company, TradingClearingRegistry tcr) { this.account = account; this.company = company; + this.tcr = tcr; } public boolean isFound() { - return error == null && account != null && company != null; + return error == null && account != null && company != null && tcr != null; } public EnumMessage getError() { @@ -134,5 +136,13 @@ public class AnltSearcher { public void setCompany(Company company) { this.company = company; } + + public TradingClearingRegistry getTcr() { + return tcr; + } + + public void setTcr(TradingClearingRegistry tcr) { + this.tcr = tcr; + } } } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf06Executor.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf06Executor.java index 541299d2d..a723241c3 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf06Executor.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf06Executor.java @@ -22,6 +22,7 @@ import ru.spcex.clearing.platform.messaging.domain.cud.clearing.AssetOperationRe import ru.spcex.clearing.platform.messaging.domain.cud.gateway.AssetOperationApprovalRequest; import ru.spcex.clearing.platform.messaging.domain.cud.gateway.SingleAssetResponse; import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; +import ru.spcex.clearing.service.integration.GatewayRequestCreator; import ru.spcex.clearing.service.registry.DmiService; import ru.spcex.clearing.service.schedule.TradingTimeService; import ru.spcex.clearing.service.validation.Sdf06NewValidationRule; @@ -35,7 +36,6 @@ import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder; import ru.spcex.platform.imdg.validation.ImdgValidationContext; import ru.spcex.platform.utils.enumeration.EnumMessage; import ru.spcex.platform.utils.enumeration.IMessageResolver; -import ru.spcex.platform.utils.number.BigDecimalUtil; import ru.spcex.platform.utils.text.TextUtil; import ru.spcex.platform.utils.validation.IValidator; import ru.spcex.platform.utils.validation.ValidatorImpl; @@ -49,6 +49,8 @@ import java.util.Map; import java.util.Optional; import java.util.function.Function; +import static ru.spcex.platform.utils.number.BigDecimalUtil.safeBD; + @Service public class Sdf06Executor { private final Logger log = LoggerFactory.getLogger(getClass()); @@ -57,6 +59,7 @@ public class Sdf06Executor { private final ImdgProvider imdgProvider; private final Imdg accImdg; private final Imdg tcrImdg; + private final Imdg cmpImdg; private final Imdg sdf06Imdg; private final Imdg sdf07Imdg; private final ImdgId idGenerator; @@ -95,6 +98,7 @@ public class Sdf06Executor { this.kafkaSender = kafkaSender; this.accImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Account, Account.class); this.tcrImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class); + this.cmpImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Company, Company.class); this.sdf06Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf06, SDf06.class); this.sdf07Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf07, SDf07.class); this.filenameObtainer = filenameObtainer; @@ -154,14 +158,15 @@ public class Sdf06Executor { statementImdg.insert(stmt); log.debug("Statement created: {}", stmt.getId()); if (InOutDirection.out.getKey().equals(stmt.getInOutDirection()) && tradingTimeService.isTradingTime()) { - requests.add(requestFromStatement(stmt, company.getTradingCode(), tcr.getCode())); + requests.add(GatewayRequestCreator.gatewayRequestPart(stmt, company.getTradingCode(), tcr.getCode())); } else { processedApproved(stmt, sDf06, now, sdf07GroupId); dmiService.setProc(stmt.getAccountId(), company.getId(), tcr.getId(), CurrencyCode.RUB.getKey(), - stmt.getAmount()); + safeBD(sDf06.getSum()), + sDf06.getNumber()); sdf07WasCreated = true; } } @@ -260,7 +265,8 @@ public class Sdf06Executor { stmt.getAddresseeId(), searchTcrOnGatewayResponse(sdf06).map(SpcexObjectBase::getId).orElse(null), CurrencyCode.RUB.getKey(), - stmt.getAmount() + safeBD(sdf06.getSum()), + sdf06.getNumber() ); } else { log.trace("statement.id={}, sdf07.id={}, sdf06.id={} rejected (by gateway answer)", @@ -318,13 +324,14 @@ public class Sdf06Executor { private Statement createStatementBySdf06(SDf06 sdf06, Long companyId, Account account) { Statement stmt = new Statement(); + stmt.setContract(sdf06.getNumber().toString()); stmt.setAddresseeId(companyId); stmt.setSenderId(Sender.Prc.getId()); stmt.setStatementType(StatementType.incr.getKey()); stmt.setComment(sdf06.getSpec()); stmt.setAccountId(account.getId()); stmt.setAccount(account.getAccount()); - BigDecimal amount = BigDecimalUtil.safeBD(sdf06.getSum()); + BigDecimal amount = safeBD(sdf06.getSum()); if (amount.compareTo(BigDecimal.ZERO) >= 0) { stmt.setInOutDirection(InOutDirection.in.getKey()); } else { @@ -354,6 +361,7 @@ public class Sdf06Executor { ImdgValidationContext ctx = new ImdgValidationContext<>(); ctx.addImdg(IMDGDistributedNames.Map_Account, accImdg); ctx.addImdg(IMDGDistributedNames.Map_TradingClearingRegistry, tcrImdg); + ctx.addImdg(IMDGDistributedNames.Map_Company, cmpImdg); ctx.setValidatedObject(sdf06); ValidatorImpl> v = new ValidatorImpl<>(ctx, Sdf06NewValidationRule.AccountPresent, Sdf06NewValidationRule.TcrPresent); 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 617e48fcc..2eaab8e12 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 @@ -18,12 +18,18 @@ 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.AssetOperationListRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.clearing.AssetOperationRequest; import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.clearing.service.AnltSearcher; import ru.spcex.clearing.service.LoggingService; import ru.spcex.clearing.service.builder.RegistryBuilder; +import ru.spcex.clearing.service.integration.GatewayRequestCreator; import ru.spcex.clearing.service.model.Result; +import ru.spcex.clearing.service.registry.DmiService; +import ru.spcex.clearing.service.schedule.TradingTimeService; import ru.spcex.clearing.service.validation.ValidationStored; import ru.spcex.clearing.session.stage.util.RegistryUtil; import ru.spcex.platform.classes.base.SpcexObjectBase; @@ -44,8 +50,10 @@ import java.time.Instant; import java.time.LocalDate; import java.time.format.DateTimeFormatter; import java.util.Collection; +import java.util.List; import java.util.Map; import java.util.Optional; +import java.util.concurrent.atomic.AtomicReference; import java.util.function.Consumer; import java.util.function.Function; import java.util.regex.Pattern; @@ -73,12 +81,16 @@ public class Sdf57Executor extends AbstractExecutor { private final Imdg currencyImdg; private final AnltSearcher anltSearcher; private final IMessageResolver messageResolver; + private final DmiService dmiService; + private final TradingTimeService timeService; + private final KafkaSender kafka; private final Pattern pattern = Pattern.compile("№.*"); public Sdf57Executor(@Qualifier("sdf57Validator") Function sDf57Validator, LoggingService errorLogger, ImdgProvider imdgProvider, - IMessageResolver errorResolver, AnltSearcher anltSearcher, IMessageResolver messageResolver) { + IMessageResolver errorResolver, AnltSearcher anltSearcher, IMessageResolver messageResolver, DmiService dmiService, TradingTimeService timeService, + @Qualifier("kafkaSenderWithoutRequestInfo") KafkaSender kafka) { this.statementImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Statement, Statement.class); this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class); this.sdf02Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf02, SDf02.class); @@ -95,6 +107,9 @@ public class Sdf57Executor extends AbstractExecutor { this.anltSearcher = anltSearcher; this.messageResolver = messageResolver; this.securityImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Security, Security.class); + this.dmiService = dmiService; + this.timeService = timeService; + this.kafka = kafka; } //todo доделать контроль sdf01 и sdf57 @@ -148,6 +163,10 @@ public class Sdf57Executor extends AbstractExecutor { if (company == null || account == null) { return Optional.empty(); } + if (AccountType.Tran.equalsByKey(account.getAccountType())) { + log.trace("sdf57.id={} account type {}, skipping statement creation", sdf57.getId(), AccountType.Tran.getKey()); + return Optional.empty(); + } Statement statement = create(sdf57, company, account, inOut); statementImdg.insert(statement); return Optional.of(statement); @@ -181,6 +200,7 @@ public class Sdf57Executor extends AbstractExecutor { .companyId(company.getId()) //раньше было stmt.getAddresseeId() .accountId(account.getId()) //раньше было stmt.getAccountId() .find(); + AtomicReference amtCached = new AtomicReference<>(); amtFound.ifPresentOrElse(rgs -> { log.debug("stmt.id={}, found AM*T.id={}, updating...", stmt.getId(), rgs.getId()); updateReg(stmt, rgs); @@ -209,6 +229,7 @@ public class Sdf57Executor extends AbstractExecutor { rgs.getBalance(), rgs.getDebit(), rgs.getCredit()); + amtCached.set(rgs); }, () -> { Registry registry = RegistryBuilder.builder(imdgProvider) .statement(stmt) @@ -228,8 +249,10 @@ public class Sdf57Executor extends AbstractExecutor { registry.getId(), registryF.getId(), registryB.getId()); + amtCached.set(registry); }); InOutDirection direction = IEnumKey.getEnumByKey(InOutDirection.class, stmt.getInOutDirection()); + boolean dm_tWasCreated = false; if (!StringUtils.isEmpty(sdf57.getSpecif()) && pattern.matcher(sdf57.getSpecif()).find() && InOutDirection.in.equals(direction)) { //todo доделать 11. Идентификация средств УК на клиринговом счете //stmt.id=340270002, sdf57.specif=Возврат депозита по договору №DT1000S001U/280623/8/7 new DM*T.id=339610032 @@ -275,6 +298,7 @@ public class Sdf57Executor extends AbstractExecutor { } else if (om_tFound.isPresent()) { Registry om_t = om_tFound.get(); Registry registryD = copyRegD(om_t, stmt, RegistryUnit.T); + dm_tWasCreated = true; registryImdg.insert(registryD); log.debug("stmt.id={}, sdf57.specif={}, found OM*T.id={} new DM*T.id={}", stmt.getId(), @@ -295,6 +319,14 @@ public class Sdf57Executor extends AbstractExecutor { stmt.setOperationStatus(OperationStatus.Executed.getKey()); statementImdg.update(stmt); log.debug("stmt.id={} status -> {}", stmt.getId(), OperationStatus.Executed.getKey()); + if (!dm_tWasCreated && timeService.isTradingTime() && InOutDirection.in.equalsByKey(stmt.getInOutDirection())) { + log.trace("sending gateway request for stmt.id={}", stmt.getId()); + kafka.sendRequestToQueue(Consts.ASSET_OPERATION, + gatewayRequest(stmt, + company.getTradingCode(), + amtCached.get().getTradingClearingRegistry())); + + } } else { log.debug("statement.id={} activness validation failed {}", stmt.getId(), messageResolver.resolve(err.get())); stmt.setErrorCodeId(err.get().getSubject().getId()); // fixme ErrorText insert @@ -316,6 +348,10 @@ public class Sdf57Executor extends AbstractExecutor { searchResult.getAccount().getId(), searchResult.getCompany().getId()); createRegistryIfNeeded.accept(new StmtCmpAcc(stmt, searchResult.getCompany(), searchResult.getAccount())); + dmiService.setOk(searchResult.getAccount().getId(), + searchResult.getCompany().getId(), + CurrencyCode.RUB.getKey(), + sdf57.getDbfId().toString()); } else { log.debug("stmt.id={} comment='{}' error: {}. Operating through DMAU registry", stmt.getId(), @@ -346,6 +382,11 @@ public class Sdf57Executor extends AbstractExecutor { }); //DMAU подразумевает увеличение остатка поля balance и credit registry.code=DMAU } + } else { + dmiService.setOk(acc.getId(), + cmp.getId(), + CurrencyCode.RUB.getKey(), + sdf57.getDbfId().toString()); } }; statementDeb.map(stmt -> new StmtCmpAcc(stmt, companyDeb, accountDeb)).ifPresent(registersUpdate); @@ -625,4 +666,10 @@ public class Sdf57Executor extends AbstractExecutor { R apply(T1 arg1, T2 arg2, T3 arg3); } + private AssetOperationListRequest gatewayRequest(Statement stmt, String tradingCode, String tcrCode) { + AssetOperationListRequest gatewayRequest = new AssetOperationListRequest(); + AssetOperationRequest req = GatewayRequestCreator.gatewayRequestPart(stmt, tradingCode, tcrCode); + gatewayRequest.setAssetOperationRequests(List.of(req)); + return gatewayRequest; + } } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/integration/GatewayRequestCreator.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/integration/GatewayRequestCreator.java new file mode 100644 index 000000000..55a10754b --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/integration/GatewayRequestCreator.java @@ -0,0 +1,18 @@ +package ru.spcex.clearing.service.integration; + +import ru.clearing.classes.statics.data.statement.Statement; +import ru.spcex.clearing.platform.messaging.domain.cud.clearing.AssetOperationRequest; +import ru.spcex.platform.enumeration.CurrencyCode; + +public class GatewayRequestCreator { + public static AssetOperationRequest gatewayRequestPart(Statement stmt, String tradingCode, String tcrCode) { + AssetOperationRequest req = new AssetOperationRequest(); + req.setStatementId(stmt.getId()); + req.setAmount(stmt.getAmount()); + req.setSecuritySymbol(CurrencyCode.RUB.getKey()); //fixme retrieve security symbol from validator + req.setTradingCode(tradingCode); + req.setCode(tcrCode); + req.setDirection(stmt.getInOutDirection()); + return req; + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/registry/DmiService.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/registry/DmiService.java index 80118bbed..aa84bcaef 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/registry/DmiService.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/registry/DmiService.java @@ -35,35 +35,43 @@ public class DmiService { Long companyId, Long tcrId, String securitySymbol, - BigDecimal balance) { - findDmi(companyId, accountId, securitySymbol).ifPresentOrElse(d__i -> { + BigDecimal balance, + BigDecimal contract) { + findDmi(companyId, accountId, securitySymbol, contract.toString()).ifPresentOrElse(d__i -> { d__i.setRegistryStatus(RegistryStatus.PROC.getKey()); d__i.setUpdated(Instant.now()); log.trace("updating D**I status to {}", d__i.getRegistryStatus()); rgsImdg.update(d__i); - }, () -> createDmi(companyId, tcrId, securitySymbol, balance, "outDocument")); + }, () -> createDmi(companyId, tcrId, securitySymbol, balance, contract.toString())); } - public void setOk(Long accountId, Long companyId, Long tcrId, String securitySymbol) { - findDmi(companyId, accountId, securitySymbol).ifPresent(d__i -> { + public void setOk(Long accountId, + Long companyId, + String securitySymbol, + String contract) { + findDmi(companyId, accountId, securitySymbol, contract).ifPresentOrElse(d__i -> { d__i.setRegistryStatus(RegistryStatus.OK.getKey()); d__i.setUpdated(Instant.now()); log.trace("updating {} status to {}", d__i.getRegistryCode(), d__i.getRegistryStatus()); rgsImdg.update(d__i); + }, () -> { + log.trace("didn't find DM*I by companyId={} accountId={} securitySymbol={} contract={}", + companyId, accountId, securitySymbol, contract); }); } - private Optional findDmi(Long companyId, Long accountId, String securitySymbol) { + private Optional findDmi(Long companyId, Long accountId, String securitySymbol, String contract) { ImdgPredicateBuilder pb = rgsImdg.predicateBuilder(); RegistryTradingParams rgsCde = RegistryTradingParams.D__I; ImdgPredicate prdct = pb.and( pb.sql(RegistryCodeSqlBuilder.getInstance(rgsCde).build()), pb.equals("companyId", companyId), pb.equals("accountId", accountId), - pb.equals("securitySymbol", securitySymbol) + pb.equals("securitySymbol", securitySymbol), + pb.equals("contract", contract) ); - log.trace("searching {} by {}", rgsCde, prdct); + log.trace("searching D**I by {}", prdct); Registry d__i = rgsImdg.getFirstObjectByPredicate(prdct); return Optional.ofNullable(d__i); } @@ -79,8 +87,8 @@ public class DmiService { securitySymbol, rgsCde); if (am_t.isEmpty()) { - log.error("{} not found for tcr.id={} company.id={} security.symbol={}", - rgsCde, tcrId, companyId, securitySymbol); + log.error("AM*T not found for tcr.id={} company.id={} security.symbol={}", + tcrId, companyId, securitySymbol); return;// Optional.empty(); } Registry rgsD = am_t.get().clone(); diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/Sdf06NewValidationRule.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/Sdf06NewValidationRule.java index 17a3afccb..a7e990f5f 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/Sdf06NewValidationRule.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/Sdf06NewValidationRule.java @@ -31,6 +31,16 @@ import java.util.Optional; import java.util.function.Function; public enum Sdf06NewValidationRule implements IValidationRule> { + Fields() { + @Override + public Optional validate(ImdgValidationContext context) { + SDf06 validatedObject = context.getValidatedObject(); + if (validatedObject.getNumber() == null) { + return of(ClearingError.WrongField, "number"); + } + return empty(); + } + }, CompanySymbolAndCompanyPresent() { @Override public Optional validate(ImdgValidationContext context) {