From 93ffe36ea6902f14887a03ab4beac48086b7b7d6 Mon Sep 17 00:00:00 2001 From: ialbert Date: Tue, 19 Sep 2023 18:16:27 +0300 Subject: [PATCH] =?UTF-8?q?SDF54/SDF55/SDF57=20-=20D**I=20=D1=80=D0=B5?= =?UTF-8?q?=D0=B3=D0=B8=D1=81=D1=82=D1=80=D1=8B?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../ru/spcex/clearing/config/KafkaConfig.java | 2 +- .../clearing/config/ValidationConfig.java | 5 +- .../clearing/service/EventsReceiver.java | 4 + .../service/executors/Sdf06Executor.java | 6 +- .../service/executors/Sdf55Executor.java | 14 ++- .../service/executors/Sdf57Executor.java | 30 +++--- .../integration/GatewayRequestCreator.java | 34 +++++-- .../PaymentInstructionOutboundService.java | 54 +++++++--- .../clearing/service/registry/DmiService.java | 98 +++++++++++++++---- .../PaymentOutboundValidationRule.java | 30 ++++++ .../validation/Sdf55ValidationRule.java | 11 ++- .../service/validation/ValidationStored.java | 2 +- .../session/stage/impl/GatewayRequester.java | 38 +++---- .../stage/impl/InspectionObligations.java | 4 +- .../InspectionObligationsDepositReturn.java | 4 +- .../src/main/resources/logback.xml | 2 +- 16 files changed, 247 insertions(+), 91 deletions(-) diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/KafkaConfig.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/KafkaConfig.java index 1c6662ff3..2726478a3 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/KafkaConfig.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/KafkaConfig.java @@ -44,7 +44,7 @@ public class KafkaConfig { gatewayConsumer.setSessionTimeoutMs(consumer.getSessionTimeoutMs()); gatewayConsumer.setAutoOffsetReset(consumer.getAutoOffsetReset()); gatewayConsumer.setEnableAutoCommit(consumer.getEnableAutoCommit()); - gatewayConsumer.setGroupId(consumer.getGroupId() + "-session-" + groupId.getAndIncrement()); + gatewayConsumer.setGroupId(consumer.getGroupId() + "-gateway-" + groupId.getAndIncrement()); return KafkaConsumerFactory.consumer(gatewayConsumer); }; } 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 c5cb3e78e..82692c21e 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 @@ -398,6 +398,7 @@ public class ValidationConfig { ctx.setValidatedObject(pmtOut); ctx.addImdg(IMDGDistributedNames.Map_Company, imdgCompany); ctx.addImdg(IMDGDistributedNames.Map_Account, imdgAccount); + ctx.addImdg(IMDGDistributedNames.Map_TradingClearingRegistry, imdgTradingClearingRegistry); return new ValidatorImpl<>(ctx, PaymentOutboundValidationRule.RequiredFields, IdPresentRule.instance("senderId", @@ -413,7 +414,9 @@ public class ValidationConfig { ClearingError.RequiredFieldEmpty, ClearingError.DictionaryNotFound), PaymentOutboundValidationRule.CreditLegAccount, - PaymentOutboundValidationRule.DebitLegAccount); + PaymentOutboundValidationRule.DebitLegAccount, + PaymentOutboundValidationRule.AddresseePresent, + PaymentOutboundValidationRule.TcrPresent); }; } 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 4ba7ddd47..cafd85d1c 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 @@ -18,6 +18,7 @@ import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonIdRequest; import ru.spcex.clearing.platform.messaging.domain.cud.gateway.AssetOperationApprovalRequest; import ru.spcex.clearing.platform.messaging.domain.cud.payment.PIClearingOutbondActionNewRequest; import ru.spcex.clearing.platform.messaging.domain.cud.registry.RegistryChangeRefundDateRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.registry.RegistryChangeStatusExtractRequest; import ru.spcex.clearing.platform.messaging.domain.cud.registry.RegistryReturnDepositRequest; import ru.spcex.clearing.platform.messaging.domain.cud.registry.RegistrySplitDepositActionRequest; import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest; @@ -155,6 +156,9 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean { callback(RegistryReturnDepositRequest.class) .setFunction(registryService::returnDeposit) .forDestination(Consts.REGISTRY_RETURN_DEPOSIT_ACTION, callbacks::put); + callback(RegistryChangeStatusExtractRequest.class) + .setFunction(registryService::changeStatusExtract) + .forDestination(Consts.REGISTRY_CHANGE_STATUS_EXTRACT_ACTION, callbacks::put); callback(RegistryChangeRefundDateRequest.class) .setFunction(registryService::changeRefundDate) .forDestination(Consts.REGISTRY_CHANGE_REFUND_DATE_ACTION, callbacks::put); 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 b07a3e798..9cb346202 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 @@ -152,11 +152,11 @@ public class Sdf06Executor { statementImdg.insert(stmt); log.debug("Statement created: {}", stmt.getId()); if (InOutDirection.out.getKey().equals(stmt.getInOutDirection()) && tradingTimeService.isTradingTime()) { - requests.add(GatewayRequestCreator.gatewayRequestPart(stmt, company.getTradingCode(), tcr.getCode())); + requests.add(GatewayRequestCreator.from(stmt, company.getTradingCode(), tcr.getCode())); } else { processedApproved(stmt, sDf06, now, sdf07GroupId); if (tcr != null) { - dmiService.setProc(tcr.getId(), + dmiService.setProcContract(tcr.getId(), CurrencyCode.RUB.getKey(), safeBD(sDf06.getSum()), sDf06.getNumber().toString()); @@ -263,7 +263,7 @@ public class Sdf06Executor { Instant updatedTime = Instant.now(); if (gatewayMsg.isApproved()) { processedApproved(stmt, sdf06, updatedTime); - dmiService.setProc(searchTcrOnGatewayResponse(sdf06).map(SpcexObjectBase::getId).orElse(null), + dmiService.setProcContract(searchTcrOnGatewayResponse(sdf06).map(SpcexObjectBase::getId).orElse(null), CurrencyCode.RUB.getKey(), safeBD(sdf06.getSum()), sdf06.getNumber().toString() diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf55Executor.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf55Executor.java index 133897611..b2e49539e 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf55Executor.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf55Executor.java @@ -50,7 +50,7 @@ public class Sdf55Executor { this.assets = assets; } - public void execute(BaseRequest systemRequest) { + public synchronized void execute(BaseRequest systemRequest) { StatementRequest payload = systemRequest.getRequestPayload(); log.debug("executing {}.generationId={}", payload.getTable(), payload.getGroupId()); IValidator validator = validation.apply(payload); @@ -68,12 +68,10 @@ public class Sdf55Executor { pmt.setTransactionStatus(TransactionStatus.ok.getKey()); pmt.setUpdated(Instant.now()); pmtImdg.update(pmt); - if (tcr != null) { - dmiService.setOk(tcr.getId(), CurrencyCode.RUB.getKey(), sDf54.getDocnm_ref()); - assets.process(assetTrio.a__b(), - assetTrio.a__t(), - assetTrio.a__f(), - BigDecimal.ZERO); - } + dmiService.setContract(tcr.getId(), CurrencyCode.RUB.getKey(), sDf54.getDocnm_ref(), sDf55.getDocnmprev()); + assets.process(assetTrio.a__b(), + assetTrio.a__t(), + assetTrio.a__f(), + BigDecimal.ZERO); } } 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 a0299ad48..27db1406e 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 @@ -25,6 +25,7 @@ import ru.spcex.clearing.platform.messaging.domain.cud.clearing.AssetOperationLi 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.AssetTrio; import ru.spcex.clearing.service.LoggingService; import ru.spcex.clearing.service.builder.RegistryBuilder; import ru.spcex.clearing.service.integration.GatewayRequestCreator; @@ -56,7 +57,6 @@ 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; @@ -197,7 +197,7 @@ public class Sdf57Executor extends AbstractExecutor { statementDeb.map(SpcexObjectBase::getId).orElse(null), statementCred.map(SpcexObjectBase::getId).orElse(null)); - Function> createRegistryIfNeeded = stmtCmpAcc -> { + Function> createRegistryIfNeeded = stmtCmpAcc -> { Statement stmt = stmtCmpAcc.statement(); Company company = stmtCmpAcc.company(); Account account = stmtCmpAcc.account(); @@ -208,8 +208,7 @@ public class Sdf57Executor extends AbstractExecutor { .companyId(company.getId()) //раньше было stmt.getAddresseeId() .accountId(account.getId()) //раньше было stmt.getAccountId() .find(); - AtomicReference amtCached = new AtomicReference<>(); - amtFound.ifPresentOrElse(rgs -> { + AssetTrio asts = amtFound.map(rgs -> { log.debug("stmt.id={}, found AM*T.id={}, updating...", stmt.getId(), rgs.getId()); updateReg(stmt, rgs); //добавил создание если не найдены @@ -239,8 +238,8 @@ public class Sdf57Executor extends AbstractExecutor { rgs.getBalance(), rgs.getDebit(), rgs.getCredit()); - amtCached.set(rgs); - }, () -> { + return new AssetTrio(registryUnitF, registryUnitB, rgs); + }).orElseGet(() -> { Registry registry = RegistryBuilder.builder(imdgProvider) .statement(stmt) .company(company) @@ -259,7 +258,7 @@ public class Sdf57Executor extends AbstractExecutor { registry.getId(), registryF.getId(), registryB.getId()); - amtCached.set(registry); + return new AssetTrio(registryF, registryB, registry); }); InOutDirection direction = IEnumKey.getEnumByKey(InOutDirection.class, stmt.getInOutDirection()); boolean dm_tWasCreated = false; @@ -334,10 +333,9 @@ public class Sdf57Executor extends AbstractExecutor { kafka.sendRequestToQueue(Consts.ASSET_OPERATION, gatewayRequest(stmt, company.getTradingCode(), - amtCached.get().getTradingClearingRegistry())); - + asts.a__t().getTradingClearingRegistry())); } - return Optional.ofNullable(amtCached.get().getTradingClearingRegistryId()); + return Optional.of(asts); } else { log.debug("statement.id={} activness validation failed {}", stmt.getId(), messageResolver.resolve(err.get())); stmt.setErrorCodeId(err.get().getSubject().getId()); // fixme ErrorText insert @@ -350,7 +348,7 @@ public class Sdf57Executor extends AbstractExecutor { Statement stmt = stmtCmpAcc.statement(); Company cmp = stmtCmpAcc.company(); Account acc = stmtCmpAcc.account(); - Optional tcrId = createRegistryIfNeeded.apply(stmtCmpAcc); + Optional asts = createRegistryIfNeeded.apply(stmtCmpAcc); if (accountIsAnlt(acc)) { AnltSearcher.AnltSearch searchResult = anltSearcher.loadByAnlt(stmt.getComment()); if (searchResult.isFound()) { @@ -359,11 +357,12 @@ public class Sdf57Executor extends AbstractExecutor { stmt.getComment(), searchResult.getAccount().getId(), searchResult.getCompany().getId()); - createRegistryIfNeeded.apply(new StmtCmpAcc(stmt, searchResult.getCompany(), searchResult.getAccount())); + asts = createRegistryIfNeeded.apply(new StmtCmpAcc(stmt, searchResult.getCompany(), searchResult.getAccount())); if (OperationStatus.Executed.equalsByKey(stmt.getOperationStatus())) { dmiService.setOk(searchResult.getTcr().getId(), CurrencyCode.RUB.getKey(), sdf57.getDbfId().toString()); + asts.ifPresent(trio -> assets.process(trio.a__b(), trio.a__t(), trio.a__f(), BigDecimal.ZERO)); } } else { log.debug("stmt.id={} comment='{}' error: {}. Operating through DMAU registry", @@ -396,10 +395,11 @@ public class Sdf57Executor extends AbstractExecutor { //DMAU подразумевает увеличение остатка поля balance и credit registry.code=DMAU } } else if (OperationStatus.Executed.equalsByKey(stmt.getOperationStatus())) { - if (tcrId.isPresent()) { - dmiService.setOk(tcrId.get(), + if (asts.isPresent()) { + dmiService.setOk(asts.get().a__t().getTradingClearingRegistryId(), CurrencyCode.RUB.getKey(), sdf57.getDbfId().toString()); + asts.ifPresent(trio -> assets.process(trio.a__b(), trio.a__t(), trio.a__f(), BigDecimal.ZERO)); } else { log.error("stmt.id={} operationStatus={} but couldn't extract TCR.id for DM*I update", stmt.getId(), @@ -679,7 +679,7 @@ public class Sdf57Executor extends AbstractExecutor { private AssetOperationListRequest gatewayRequest(Statement stmt, String tradingCode, String tcrCode) { AssetOperationListRequest gatewayRequest = new AssetOperationListRequest(); - AssetOperationRequest req = GatewayRequestCreator.gatewayRequestPart(stmt, tradingCode, tcrCode); + AssetOperationRequest req = GatewayRequestCreator.from(stmt, tradingCode, tcrCode); Session existActiveSession = sessionImdg.getFirstObjectByFieldValues(Map.of( "workflowStatus", SessionStatus.ACTV.getKey() )); 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 index 313331997..885c2e219 100644 --- 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 @@ -1,5 +1,7 @@ package ru.spcex.clearing.service.integration; +import ru.clearing.classes.statics.data.payment.PaymentInstruction; +import ru.clearing.classes.statics.data.registry.Registry; import ru.clearing.classes.statics.data.statement.Statement; import ru.spcex.clearing.platform.messaging.domain.cud.clearing.AssetOperationRequest; import ru.spcex.platform.enumeration.CurrencyCode; @@ -8,7 +10,7 @@ import ru.spcex.platform.enumeration.InOutDirection; import java.math.BigDecimal; public class GatewayRequestCreator { - public static AssetOperationRequest gatewayRequestPart(Statement stmt, String tradingCode, String tcrCode) { + public static AssetOperationRequest from(Statement stmt, String tradingCode, String tcrCode) { AssetOperationRequest req = new AssetOperationRequest(); req.setEntityId(stmt.getId()); req.setAmount(stmt.getAmount()); @@ -19,12 +21,30 @@ public class GatewayRequestCreator { return req; } - public static AssetOperationRequest gatewayRequestPart(Long entityId, - Long sessionId, - BigDecimal amount, - InOutDirection direction, - String tradingCode, - String tcrCode) { + public static AssetOperationRequest from(Registry om_t) { + return from(om_t.getId(), + null, + om_t.getBalance(), + InOutDirection.out, + om_t.getTradingCode(), + om_t.getTradingClearingRegistry()); + } + + public static AssetOperationRequest from(PaymentInstruction pmt, String tradingCode, String tcr) { + return from(pmt.getId(), + null, + pmt.getCreditLeg_amount(), + InOutDirection.out, + tradingCode, + tcr); + } + + public static AssetOperationRequest from(Long entityId, + Long sessionId, + BigDecimal amount, + InOutDirection direction, + String tradingCode, + String tcrCode) { AssetOperationRequest req = new AssetOperationRequest(); req.setSessionId(sessionId); req.setEntityId(entityId); diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/payment/PaymentInstructionOutboundService.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/payment/PaymentInstructionOutboundService.java index 86ddcc805..4aabdeb6c 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/payment/PaymentInstructionOutboundService.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/payment/PaymentInstructionOutboundService.java @@ -14,8 +14,10 @@ import ru.clearing.classes.statics.data.registry.TradingClearingRegistry; import ru.clearing.classes.statics.data.sdf.SDf54; import ru.spcex.clearing.error.ClearingError; import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.notification.NotificationSender; import ru.spcex.clearing.platform.messaging.domain.BaseRequest; import ru.spcex.clearing.platform.messaging.domain.Consts; +import ru.spcex.clearing.platform.messaging.domain.cud.clearing.AssetOperationRequest; import ru.spcex.clearing.platform.messaging.domain.cud.clearing.SdfClearingRequest; import ru.spcex.clearing.platform.messaging.domain.cud.payment.PIClearingOutbondActionNewRequest; import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate; @@ -24,14 +26,15 @@ import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.clearing.service.AssetTrio; import ru.spcex.clearing.service.Sdf54Creator; import ru.spcex.clearing.service.builder.PaymentInstructionBuilderV2; +import ru.spcex.clearing.service.integration.GatewayRequestCreator; import ru.spcex.clearing.service.registry.AssetTBFProcessing; import ru.spcex.clearing.service.registry.DmiService; import ru.spcex.clearing.service.registry.RegistryManager; +import ru.spcex.clearing.service.schedule.TradingTimeService; import ru.spcex.clearing.service.validation.ValidationStored; +import ru.spcex.clearing.session.stage.impl.GatewayRequester; import ru.spcex.clearing.util.security.UserRoleVerification; -import ru.spcex.platform.enumeration.CurrencyCode; -import ru.spcex.platform.enumeration.TransactionStatus; -import ru.spcex.platform.enumeration.UserRole; +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; @@ -44,6 +47,7 @@ import java.math.BigDecimal; import java.time.Instant; import java.util.Optional; import java.util.function.Function; +import java.util.function.Supplier; import static ru.spcex.platform.utils.number.BigDecimalUtil.safeBD; @@ -67,12 +71,15 @@ public class PaymentInstructionOutboundService { private final KafkaSender kafkaSender; private final AssetTBFProcessing assetsMng; private final DmiService dmiService; + private final TradingTimeService time; + private final GatewayRequester gateway; + private final NotificationSender notification; @Autowired public PaymentInstructionOutboundService(UserRoleVerification rights, IMessageResolver msgs, ImdgProvider imdgProvider, - Function validation, RegistryManager rgsMng, KafkaSender kafkaSender, AssetTBFProcessing assetsMng, DmiService dmiService) { + Function validation, RegistryManager rgsMng, KafkaSender kafkaSender, AssetTBFProcessing assetsMng, DmiService dmiService, TradingTimeService time, GatewayRequester gateway, NotificationSender notification) { this.imdgProvider = imdgProvider; this.rights = rights; this.msgs = msgs; @@ -90,6 +97,10 @@ public class PaymentInstructionOutboundService { this.kafkaSender = kafkaSender; this.assetsMng = assetsMng; this.dmiService = dmiService; + this.time = time; + this.gateway = gateway; + this.gateway.setName("PaymentInstructionOutboundService|PI"); + this.notification = notification; } public RequestInfoUpdate sendOutBoundPayment(BaseRequest req) { @@ -106,14 +117,13 @@ public class PaymentInstructionOutboundService { } Account accCred = validator.getStored(ValidationStored.PIOutboundAccountCred); Account accDeb = validator.getStored(ValidationStored.PIOutboundAccountDeb); - Company addressee = cmpImdg.getSingleObjectByID(payload.getAddresseeId()); + Company addressee = validator.getStored(ValidationStored.PIOOutboundAddressee); + TradingClearingRegistry tcr = validator.getStored(ValidationStored.PIOutboundTcr); Company sender = cmpImdg.getSingleObjectByID(payload.getSenderId()); BigDecimal amount = safeBD(payload.getCreditLeg_amount()); log.debug("all checks passed, accCred.id={}, accDeb.id={}, addressee.id={}, sender.id={}, amount: {}", accCred.getId(), accDeb.getId(), addressee.getId(), sender.getId(), amount); - TradingClearingRegistry tcr = tcrImdg.getFirstObjectBySQL("moneyAccountId = '%d' and companyId = %d" - .formatted(accCred.getId(), addressee.getId())); String purpose = tcr != null ? "Вывод средств по ТКР " + tcr.getCode() + "." : "Вывод средств."; if (payload.getPaymentPurpose() != null) { purpose += " " + payload.getPaymentPurpose(); @@ -141,7 +151,7 @@ public class PaymentInstructionOutboundService { pmt.setCreditLeg_securityId(currency.getId()); pmt.setDebitLeg_securityId(currency.getId()); } - pmt.setTransactionStatus(TransactionStatus.exec.getKey()); + pmt.setTransactionStatus(TransactionStatus.stld.getKey()); //fixme pmt.setDocumentNumber(registry.getSecurityId()); pmtImdg.insert(pmt); @@ -162,12 +172,28 @@ public class PaymentInstructionOutboundService { sDf54.setGenerationId(pmt.getId()); sdf54Imdg.insert(sDf54); log.debug("new sdf54.id: {}", sDf54.getId()); - if (tcr != null) { - dmiService.setProc(tcr.getId(), - CurrencyCode.RUB.getKey(), - safeBD(amount).negate(), - sDf54.getDocnm_ref()); - assetsMng.process(assets.get().a__b(), assets.get().a__t(), assets.get().a__f(), BigDecimal.ZERO); + dmiService.setProcComment(tcr.getId(), + CurrencyCode.RUB.getKey(), + safeBD(amount).negate(), + sDf54.getDocnm_ref()); + assetsMng.process(assets.get().a__b(), assets.get().a__t(), assets.get().a__f(), BigDecimal.ZERO); + if (time.isTradingTime()) { + Supplier builder = () -> GatewayRequestCreator.from(pmt, + assets.get().a__b().getTradingCode(), + tcr.getCode()); + Optional gatewayOk = gateway.gatewayRequestAndWait(builder); + if (gatewayOk.isEmpty() || gatewayOk.get().equals(Boolean.FALSE)) { + log.error("sdf54.id={} pmt.id={} gateway {}", sDf54.getId(), + pmt.getId(), + gatewayOk.isPresent() ? "failed" : "timeout"); + notification.sendNotification(ObjectType.gateway, + gatewayOk.isPresent() ? "gateway not approved" : "gateway timeout", + Priority.HIGH); + pmt.setTransactionStatus(TransactionStatus.fail.getKey()); + pmt.setUpdated(now); + pmtImdg.update(pmt); + return error(req.getId(), ClearingError.GeneralError); + } } SdfClearingRequest exp = new SdfClearingRequest(); exp.setGroupId(sDf54.getGenerationId()); 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 9ae740916..41e93d7ec 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 @@ -31,23 +31,76 @@ public class DmiService { this.rgsMng = rgsMng; } - public void setProc(Long tcrId, - String securitySymbol, - BigDecimal balance, - String contract) { - findDmi(tcrId, securitySymbol, contract.toString()).ifPresentOrElse(d__i -> { - d__i.setRegistryStatus(RegistryStatus.PROC.getKey()); + /** + * ищем D**I по tcrId/security/contract + * обновляем либо создаем + */ + public void setProcContract(Long tcrId, + String securitySymbol, + BigDecimal balance, + String contract) { + findDmiByContract(tcrId, securitySymbol, contract).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(tcrId, securitySymbol, balance) + .ifPresent(d__i -> { + d__i.setContract(contract); + rgsImdg.insert(d__i); + log.debug("created {}.id={}", d__i.getRegistryCode(), d__i.getId()); + }) + ); + } + + /** + * ищем D**I по tcrId/security/comment + * обновляем либо создаем + */ + public void setProcComment(Long tcrId, + String securitySymbol, + BigDecimal balance, + String comment) { + findDmiByComment(tcrId, securitySymbol, comment).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(tcrId, securitySymbol, balance) + .ifPresent(d__i -> { + d__i.setComment(comment); + rgsImdg.insert(d__i); + log.debug("created {}.id={}", d__i.getRegistryCode(), d__i.getId()); + }) + ); + } + + /** + * ищем D**I по tcrId/security/comment + * добавляем поле contract + */ + public void setContract(Long tcrId, + String securitySymbol, + String comment, + String contract) { + findDmiByComment(tcrId, securitySymbol, comment).ifPresentOrElse(d__i -> { + d__i.setContract(contract); d__i.setUpdated(Instant.now()); - log.trace("updating D**I status to {}", d__i.getRegistryStatus()); rgsImdg.update(d__i); - }, () -> createDmi(tcrId, securitySymbol, balance, contract.toString())); + log.trace("set {}.id={} contract {}", d__i.getRegistryCode(), d__i.getId(), contract); + }, () -> log.trace("didn't find DM*I by tcr.id={} securitySymbol={} comment={}", + tcrId, securitySymbol, comment)); } + /** + * ищем D**I по tcrId/security/contract + * задаем статус ОК + */ public void setOk(Long tcrId, String securitySymbol, String contract) { - findDmi(tcrId, securitySymbol, contract).ifPresentOrElse(d__i -> { + findDmiByContract(tcrId, 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()); @@ -58,7 +111,7 @@ public class DmiService { }); } - private Optional findDmi(Long tcrId, String securitySymbol, String contract) { + private Optional findDmiByContract(Long tcrId, String securitySymbol, String contract) { ImdgPredicateBuilder pb = rgsImdg.predicateBuilder(); RegistryTradingParams rgsCde = RegistryTradingParams.D__I; ImdgPredicate prdct = pb.and( @@ -72,10 +125,23 @@ public class DmiService { return Optional.ofNullable(d__i); } - private void createDmi(Long tcrId, + private Optional findDmiByComment(Long tcrId, String securitySymbol, String comment) { + ImdgPredicateBuilder pb = rgsImdg.predicateBuilder(); + RegistryTradingParams rgsCde = RegistryTradingParams.D__I; + ImdgPredicate prdct = pb.and( + pb.sql(RegistryCodeSqlBuilder.getInstance(rgsCde).build()), + pb.equals("tradingClearingRegistryId", tcrId), + pb.equals("securitySymbol", securitySymbol), + pb.equals("comment", comment) + ); + log.trace("searching D**I by {}", prdct); + Registry d__i = rgsImdg.getFirstObjectByPredicate(prdct); + return Optional.ofNullable(d__i); + } + + private Optional createDmi(Long tcrId, String securitySymbol, - BigDecimal summ, - String outDocument) { + BigDecimal summ) { RegistryTradingParams rgsCde = RegistryTradingParams.AM_T; ImdgPredicateBuilder pb = rgsImdg.predicateBuilder(); ImdgPredicate prdct = pb.and( @@ -87,7 +153,7 @@ public class DmiService { if (am_t == null) { log.error("AM*T not found for tcr.id={} security.symbol={}", tcrId, securitySymbol); - return; + return Optional.empty(); } Registry rgsD = am_t.clone(); rgsD.setRegistryDesignation(RegistryDesignation.D.getKey()); @@ -100,8 +166,6 @@ public class DmiService { rgsD.setDiffBalance(BigDecimal.ZERO); rgsD.setCheckBalance(BigDecimal.ZERO); rgsD.setRegistryStatus(RegistryStatus.PROC.getKey()); - rgsD.setContract(outDocument); - rgsImdg.insert(rgsD); - log.debug("created {}.id={}", rgsD.getRegistryCode(), rgsD.getId()); + return Optional.of(rgsD); } } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/PaymentOutboundValidationRule.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/PaymentOutboundValidationRule.java index c5425418a..468308dad 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/PaymentOutboundValidationRule.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/PaymentOutboundValidationRule.java @@ -3,6 +3,8 @@ package ru.spcex.clearing.service.validation; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import ru.clearing.classes.statics.data.account.Account; +import ru.clearing.classes.statics.data.company.Company; +import ru.clearing.classes.statics.data.registry.TradingClearingRegistry; import ru.spcex.clearing.error.ClearingError; import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.platform.messaging.domain.cud.payment.PIClearingOutbondActionNewRequest; @@ -80,6 +82,34 @@ public enum PaymentOutboundValidationRule implements IValidationRule validate(ImdgValidationContext context) { + PIClearingOutbondActionNewRequest validatedObject = context.getValidatedObject(); + Imdg cmpImdg = context.obtainMap(IMDGDistributedNames.Map_Company, Company.class); + Company addressee = cmpImdg.getSingleObjectByID(validatedObject.getAddresseeId()); + if (addressee == null) { + return of(ClearingError.WrongFieldValue, "addresseeId"); + } + context.storeObject(ValidationStored.PIOOutboundAddressee, addressee); + return Optional.empty(); + } + }, + TcrPresent() { + @Override + public Optional validate(ImdgValidationContext context) { + Account accCred = context.getStoredObject(ValidationStored.PIOutboundAccountCred); + Company addressee = context.getStoredObject(ValidationStored.PIOOutboundAddressee); + Imdg tcrImdg = context.obtainMap(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class); + TradingClearingRegistry tcr = tcrImdg.getFirstObjectBySQL("moneyAccountId = '%d' and companyId = %d" + .formatted(accCred.getId(), addressee.getId())); + if (tcr == null) { + return of(ClearingError.RecordNotFound, "TCR"); + } + context.storeObject(ValidationStored.PIOutboundTcr, tcr); + return Optional.empty(); + } + } ; private final static Logger log = LoggerFactory.getLogger(PaymentOutboundValidationRule.class); @Override diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/Sdf55ValidationRule.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/Sdf55ValidationRule.java index c81697f3f..60c584691 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/Sdf55ValidationRule.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/Sdf55ValidationRule.java @@ -77,13 +77,18 @@ public enum Sdf55ValidationRule implements IValidationRule gatewayRequestAndWait(Registry om_t) { + public Optional gatewayRequestAndWait(Supplier reqBuilder) { SingleAssetResponse responseReceived = null; + Long entityId = null; gatewayLock.lock(); try { if (this.omtId != null) { - log.error("om*t.id={} already sent to gateway", this.omtId); + log.error("{}.id={} already sent to gateway", getName(), this.omtId); throw new IllegalStateException(" forbid reusing GatewayRequester instance in multiple threads"); } //send Request AssetOperationListRequest assetOperationListRequest = new AssetOperationListRequest(); Collection requests = new ArrayList<>(); - AssetOperationRequest req = GatewayRequestCreator.gatewayRequestPart( - om_t.getId(), - null, - om_t.getBalance(), - InOutDirection.out, - om_t.getTradingCode(), - om_t.getTradingClearingRegistry()); + AssetOperationRequest req = reqBuilder.get(); + entityId = req.getEntityId(); requests.add(req); assetOperationListRequest.setAssetOperationRequests(requests); - log.debug("sending om*t.id={} to gateway", om_t.getId()); + log.debug("sending {}.id={} to gateway", getName(), entityId); kafkaSender.sendRequestToQueue(Consts.ASSET_OPERATION, assetOperationListRequest); - this.omtId = om_t.getId(); + this.omtId = req.getEntityId(); boolean gatewayReceived = gatewayCondition.await(GATEWAY_TIMEOUT, TimeUnit.SECONDS); if (gatewayReceived) { responseReceived = assetOperationApprovalRequest; @@ -103,10 +97,10 @@ public class GatewayRequester extends QueueConsumer implements InitializingBean } if (responseReceived != null) { boolean approved = responseReceived.isApproved(); - log.debug("om*t.id={} answer from gateway received. approval: {}", om_t.getId(), approved); + log.debug("{}.id={} answer from gateway received. approval: {}", getName(), entityId, approved); return Optional.of(approved); } else { - log.error("om*t.id={} answer from gateway not received.", om_t.getId()); + log.error("{}.id={} answer from gateway not received.", getName(), entityId); return Optional.empty(); } } @@ -122,13 +116,13 @@ public class GatewayRequester extends QueueConsumer implements InitializingBean gatewayLock.lock(); try { if (omtId == null) { - log.trace("gateway BaseRequest.id={} currently not waiting for OM*T answer", sysReq.getId()); + log.trace("gateway BaseRequest.id={} currently not waiting for {} answer", sysReq.getId(), getName()); return; } Optional assetResponseFound = approvals.stream().filter(a -> omtId.equals(a.getEntityId())).findFirst(); if (assetResponseFound.isPresent()) { - log.debug("gateway BaseRequest.id={} OM*T answer found for OM*T.id={}", sysReq.getId(), omtId); + log.debug("gateway BaseRequest.id={} answer found for {}.id={}", sysReq.getId(), getName(), omtId); omtId = null; assetOperationApprovalRequest = assetResponseFound.get(); gatewayCondition.signal(); @@ -137,4 +131,12 @@ public class GatewayRequester extends QueueConsumer implements InitializingBean gatewayLock.unlock(); } } + + public String getName() { + return name != null ? name : "entity"; + } + + public void setName(String name) { + this.name = name; + } } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InspectionObligations.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InspectionObligations.java index 55952100b..327b6d8b2 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InspectionObligations.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InspectionObligations.java @@ -12,6 +12,7 @@ import ru.clearing.classes.statics.data.registry.Registry; import ru.spcex.clearing.error.ClearingError; import ru.spcex.clearing.error.RgsError; import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.service.integration.GatewayRequestCreator; import ru.spcex.clearing.service.registry.AssetTBFProcessing; import ru.spcex.clearing.service.registry.RegistryManager; import ru.spcex.clearing.service.schedule.TradingTimeService; @@ -64,6 +65,7 @@ public class InspectionObligations implements ISessionStage { this.registryManager = registryManager; this.assets = assets; this.gateway = gateway; + this.gateway.setName("InspectionObligations|OM*T"); this.tradingTimeService = new TradingTimeService(imdgProvider); } @@ -133,7 +135,7 @@ public class InspectionObligations implements ISessionStage { Optional omt = group.stream().filter(rgs -> RegistryManager.equalsByCode(OM_T, rgs)).findFirst(); if (SessionType.FINL.equals(sessionType) && tradingTimeService.isTradingTime() && omt.isPresent()) { - Optional gatewayReceived = gateway.gatewayRequestAndWait(omt.get()); + Optional gatewayReceived = gateway.gatewayRequestAndWait(() -> GatewayRequestCreator.from(omt.get())); if (gatewayReceived.isEmpty()) { log.error("Gateway not received response for groupId: {}. ", entry.getKey()); } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InspectionObligationsDepositReturn.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InspectionObligationsDepositReturn.java index 203fc478d..0e9742fb9 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InspectionObligationsDepositReturn.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InspectionObligationsDepositReturn.java @@ -9,6 +9,7 @@ import org.springframework.stereotype.Service; import ru.clearing.classes.statics.data.company.ClearingMemberCategory; import ru.clearing.classes.statics.data.registry.Registry; import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.service.integration.GatewayRequestCreator; import ru.spcex.clearing.service.registry.AssetTBFProcessing; import ru.spcex.clearing.service.registry.RegistryManager; import ru.spcex.clearing.service.schedule.TradingTimeService; @@ -52,6 +53,7 @@ public class InspectionObligationsDepositReturn implements ISessionStage { this.assets = assets; this.tradingTimeService = tradingTimeService; this.gateway = gateway; + this.gateway.setName("InspectionObligationsDepositReturn|OM*T"); } @Override @@ -66,7 +68,7 @@ public class InspectionObligationsDepositReturn implements ISessionStage { private boolean gatewayIfNeeded(Registry omt) { if (SessionType.FINL.equals(sessionType) && tradingTimeService.isTradingTime()) { - Optional gatewayReceived = gateway.gatewayRequestAndWait(omt); + Optional gatewayReceived = gateway.gatewayRequestAndWait(() -> GatewayRequestCreator.from(omt)); if (gatewayReceived.isEmpty()) { log.error("gateway not received response for om*t.id: {}. ", omt.getId()); } else { diff --git a/clearing-parent/clearing-service/src/main/resources/logback.xml b/clearing-parent/clearing-service/src/main/resources/logback.xml index ca7b0bb5b..c21ae570e 100644 --- a/clearing-parent/clearing-service/src/main/resources/logback.xml +++ b/clearing-parent/clearing-service/src/main/resources/logback.xml @@ -45,7 +45,7 @@ - +