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 f1bea5dae..1867d0160 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 @@ -4,6 +4,7 @@ import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import ru.clearing.classes.statics.data.account.Account; import ru.clearing.classes.statics.data.account.AccountBalance; +import ru.clearing.classes.statics.data.account.DepoAccount; import ru.clearing.classes.statics.data.company.Company; import ru.clearing.classes.statics.data.company.CompanySymbols; import ru.clearing.classes.statics.data.company.relation.Relation; @@ -46,6 +47,7 @@ public class ValidationConfig { Imdg imdgCompanySymbols; Imdg imdgRegistry; Imdg imdgAccount; + Imdg imdgDepoAccount; Imdg imdgAccountBalance; Imdg imdgSecurity; Imdg imdgMoneyMarketSecurity; @@ -70,6 +72,7 @@ public class ValidationConfig { this.fixedIncomeSecurityImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_FixedIncomeSecurity, FixedIncomeSecurity.class); this.equitySecurityImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_EquitySecurity, EquitySecurity.class); this.imdgStatement = imdgProvider.getImdg(IMDGDistributedNames.Map_Statement, Statement.class); + this.imdgDepoAccount = imdgProvider.getImdg(IMDGDistributedNames.Map_DepoAccount, DepoAccount.class); this.imdgProvider = imdgProvider; } @@ -272,6 +275,11 @@ public class ValidationConfig { ImdgValidationContext context = new ImdgValidationContext<>(); context.setValidatedObject(sDf10); context.addImdg(IMDGDistributedNames.Map_Account, imdgAccount); + context.addImdg(IMDGDistributedNames.Map_DepoAccount, imdgDepoAccount); + context.addImdg(IMDGDistributedNames.Map_TradingClearingRegistry, imdgTradingClearingRegistry); + context.addImdg(IMDGDistributedNames.Map_FixedIncomeSecurity , fixedIncomeSecurityImdg); + context.addImdg(IMDGDistributedNames.Map_MoneyMarketSecurity , imdgMoneyMarketSecurity); + context.addImdg(IMDGDistributedNames.Map_EquitySecurity , equitySecurityImdg); context.addImdg(IMDGDistributedNames.Map_CompanySymbols, imdgCompanySymbols); context.addImdg(IMDGDistributedNames.Map_Company, imdgCompany); context.addImdg(IMDGDistributedNames.Map_Registry, imdgRegistry); @@ -281,6 +289,7 @@ public class ValidationConfig { Sdf10ValidationRule.DepoAccountPresent, CompanyByTradingCodeValidationRule.instance(SDf10::getDepoCode), Sdf10ValidationRule.CompanyStatus, + Sdf10ValidationRule.TcrPresent, SecurityBySecurityCodeValidationRule.instance(SDf10::getSecurityCode), Sdf10ValidationRule.Balance ); 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 1e794c944..8c956a2b4 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 @@ -19,6 +19,7 @@ import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandR import ru.spcex.clearing.platform.messaging.domain.cud.utilities.STradesImportedRequest; import ru.spcex.clearing.platform.messaging.service.QueueConsumer; import ru.spcex.clearing.service.executors.Sdf06Executor; +import ru.spcex.clearing.service.executors.Sdf10Executor; import ru.spcex.clearing.session.stage.*; import ru.spcex.clearing.session.stage.impl.BalanceRevise; import ru.spcex.platform.enumeration.Task; @@ -39,6 +40,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean { private final ReturnDepositSession returnDepositSession; private final SessionManager sessionManager; private final Sdf06Executor sdf06Executor; + private final Sdf10Executor sdf10Executor; private final BalanceRevise balanceRevise; private final Sdf05Sender sdf05Sender; @@ -50,7 +52,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean { SecondaryAuctionT0Session secondaryAuctionT0Session, PrimaryAuctionB0Session primaryAuctionB0Session, PrimaryAuctionT0Session primaryAuctionT0Session, IntermediateMkrSession intermediateMkrSession, FinalMkrSession finalMkrSession, ReturnDepositSession returnDepositSession, SessionManager sessionManager, Sdf06Executor sdf06Executor, - BalanceRevise balanceRevise, Sdf05Sender sdf05Sender) { + Sdf10Executor sdf10Executor, BalanceRevise balanceRevise, Sdf05Sender sdf05Sender) { super(kafkaQueue, kafkaResponseQueue); this.clearingService = clearingService; this.registryService = registryService; @@ -63,6 +65,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean { this.returnDepositSession = returnDepositSession; this.sessionManager = sessionManager; this.sdf06Executor = sdf06Executor; + this.sdf10Executor = sdf10Executor; this.balanceRevise = balanceRevise; this.sdf05Sender = sdf05Sender; } @@ -102,11 +105,14 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean { .forDestination(Consts.CONTINUE_SESSION_BN_FIRST_PART, callbacks::put); callback(StatementRequest.class) - .setConsumer(sdf06Executor::createStatementAndSendGatewayCommand) + .setConsumer(sdf06Executor::execute) .forDestination(Consts.STATEMENT_PROCESS_SDF06, callbacks::put); callback(AssetOperationApprovalRequest.class) - .setConsumer(sdf06Executor::processGatewayResponse) + .setConsumer(req -> { + sdf06Executor.processGatewayResponse(req); + sdf10Executor.processGatewayResponse(req); + }) .forDestination(Consts.ASSET_OPERATION_APPROVAL, callbacks::put); callback(STradesImportedRequest.class) 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 5eac64edd..6af94fa15 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 @@ -31,6 +31,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.ImdgPredicateBuilder; import ru.spcex.platform.utils.enumeration.EnumMessage; import ru.spcex.platform.utils.enumeration.IMessageResolver; import ru.spcex.platform.utils.number.BigDecimalUtil; @@ -102,7 +103,7 @@ public class Sdf06Executor { this.filenameObtainer = filenameObtainer; } - public void createStatementAndSendGatewayCommand(BaseRequest systemRequest) { + public void execute(BaseRequest systemRequest) { Imdg sdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf06, SDf06.class); Long groupId = systemRequest.getRequestPayload().getGroupId(); Collection sdfs = sdfImdg.getCollectionObjectsByFieldValues(Map.of( @@ -209,15 +210,23 @@ public class Sdf06Executor { //может прийти неограниченно позже 06го, после того как прошел клиринг например... Instant now = Instant.now(); String fileName = null; + boolean requestIsIntendedForSdf06 = false; for (SingleAssetResponse gatewayMsg : req.getRequestPayload().getApprovals()) { //получаем запрос для текущей группы sdf06 //находим группу Long statementId = gatewayMsg.getStatementId(); - Statement stmt = statementImdg.getSingleObjectByID(statementId); + ImdgPredicateBuilder pb = statementImdg.predicateBuilder(); + Statement stmt = statementImdg.getFirstObjectByPredicate( + pb.and( + pb.equals("id", statementId), + pb.equals("inOutSDfType", InOutSDfType.type6.getKey()) + ) + ); if (stmt == null) { - log.error("Statement.id {} not found", statementId); + log.debug("Statement.id {} with type {} not found", statementId, InOutSDfType.type6.getKey()); return; } + requestIsIntendedForSdf06 = true; Long sdf06Id = stmt.getInSDfId(); SDf06 sdf06 = sdf06Imdg.getSingleObjectByID(sdf06Id); @@ -256,9 +265,13 @@ public class Sdf06Executor { statementImdg.update(stmt); } } - sendToExporter(sdf07GroupId, fileName); - log.debug("All gateway responses received for SDF06 groupId {}. SDF07 groupId {}", sdf06GroupId, sdf07GroupId); - clearContext(); + if (requestIsIntendedForSdf06) { + sendToExporter(sdf07GroupId, fileName); + log.debug("All gateway responses received for SDF06 groupId {}. SDF07 groupId {}", sdf06GroupId, sdf07GroupId); + clearContext(); + } else { + log.debug("gateway response was not for SDF06 executor"); + } } private void clearContext() { diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf10Executor.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf10Executor.java index 0027eacf2..c42913499 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf10Executor.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf10Executor.java @@ -31,6 +31,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.ImdgPredicateBuilder; import ru.spcex.platform.utils.enumeration.EnumMessage; import ru.spcex.platform.utils.enumeration.IMessageResolver; import ru.spcex.platform.utils.number.BigDecimalUtil; @@ -155,7 +156,7 @@ public class Sdf10Executor { processedApproved(stmt, sDf10, now, sdf11GroupId); sdf11WasCreated = true; } else { - requests.add(requestFromStatement(stmt, company.getTradingCode(), tcr.getCode())); + requests.add(requestFromStatement(stmt, company.getTradingCode(), tcr.getCode(), sDf10.getSecurityCode())); } } if (requests.size() > 0) { @@ -189,15 +190,23 @@ public class Sdf10Executor { public void processGatewayResponse(BaseRequest req) { Instant now = Instant.now(); + boolean requestIsIntendedForSdf10 = false; for (SingleAssetResponse gatewayMsg : req.getRequestPayload().getApprovals()) { //получаем запрос для текущей группы sdf10 //находим группу Long statementId = gatewayMsg.getStatementId(); - Statement stmt = statementImdg.getSingleObjectByID(statementId); + ImdgPredicateBuilder pb = statementImdg.predicateBuilder(); + Statement stmt = statementImdg.getFirstObjectByPredicate( + pb.and( + pb.equals("id", statementId), + pb.equals("inOutSDfType", InOutSDfType.type10.getKey()) + ) + ); if (stmt == null) { - log.error("Statement.id {} not found", statementId); + log.debug("Statement.id {} with type {} not found", statementId, InOutSDfType.type6.getKey()); return; } + requestIsIntendedForSdf10 = true; Long sdf10Id = stmt.getInSDfId(); SDf10 sdf10 = sdf10Imdg.getSingleObjectByID(sdf10Id); @@ -233,9 +242,13 @@ public class Sdf10Executor { statementImdg.update(stmt); } } - sendToExporter(sdf11GroupId); - log.debug("All gateway responses received for SDF10 groupId {}. SDF11 groupId {}", sdf10GroupId, sdf11GroupId); - clearContext(); + if (requestIsIntendedForSdf10) { + sendToExporter(sdf11GroupId); + log.debug("All gateway responses received for SDF10 groupId {}. SDF11 groupId {}", sdf10GroupId, sdf11GroupId); + clearContext(); + } else { + log.debug("gateway response was not for SDF10 executor"); + } } private void clearContext() { @@ -280,17 +293,18 @@ public class Sdf10Executor { stmt.setAmount(amount.abs()); stmt.setOperationStatus(OperationStatus.Pending.getKey()); stmt.setInSDfId(sdf10.getId()); - stmt.setInOutSDfType(InOutSDfType.type6.getKey()); + stmt.setInOutSDfType(InOutSDfType.type10.getKey()); stmt.setClearingDate(LocalDate.now()); stmt.setCreated(Instant.now()); return stmt; } - private AssetOperationRequest requestFromStatement(Statement stmt, String tradingCode, String tcrCode) { + private AssetOperationRequest requestFromStatement(Statement stmt, String tradingCode, String tcrCode, String securityCode) { AssetOperationRequest req = new AssetOperationRequest(); req.setStatementId(stmt.getId()); - req.setAmount(stmt.getAmount()); - req.setSecuritySymbol(CurrencyCode.RUB.getKey()); //fixme retrieve security symbol from validator +// req.setAmount(stmt.getAmount()); + req.setQuantity(stmt.getAmount()); + req.setSecuritySymbol(securityCode); //fixme retrieve security symbol from validator req.setTradingCode(tradingCode); req.setCode(tcrCode); req.setDirection(stmt.getInOutDirection()); diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/Sdf10ValidationRule.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/Sdf10ValidationRule.java index a75b41f15..27fa9469b 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/Sdf10ValidationRule.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/Sdf10ValidationRule.java @@ -36,7 +36,7 @@ public enum Sdf10ValidationRule implements IValidationRule accountImdg = context.obtainMap(IMDGDistributedNames.Map_DepoAccount, Account.class); + Imdg accountImdg = context.obtainMap(IMDGDistributedNames.Map_Account, Account.class); Imdg depoAccountImdg = context.obtainMap(IMDGDistributedNames.Map_DepoAccount, DepoAccount.class); Collection accs = accountImdg.getCollectionObjectsBySQL("account = '%s'".formatted(sdf10.getDepoCode())); if (accs.size() != 1) { @@ -68,7 +68,7 @@ public enum Sdf10ValidationRule implements IValidationRule validate(ImdgValidationContext context) { - DepoAccount account = context.getStoredObject(ValidationStored.Sdf10Account); + Account account = context.getStoredObject(ValidationStored.Sdf10Account); if (account == null) { return of(ClearingError.TCRegistryNotFound, ""); } @@ -92,7 +92,7 @@ public enum Sdf10ValidationRule implements IValidationRule"); @@ -114,6 +114,7 @@ public enum Sdf10ValidationRule implements IValidationRule prdctByRegistryCode = code -> prdctBldr.and( prdctBldr.equals("account", account.getAccount()), + prdctBldr.equals("securitySymbol", sdf10.getSecurityCode()), prdctBldr.equals("companyId", company.getId()), prdctBldr.sql(RegistryCodeSqlBuilder.getInstance(code).build())); Registry rgsAsf = registryImdg.getSingleObjectByPredicate(prdctByRegistryCode.apply(RegistryTradingParams.AS_F)); diff --git a/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/service/GatewayService.java b/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/service/GatewayService.java index 926a7d956..04ae7f8c3 100644 --- a/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/service/GatewayService.java +++ b/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/service/GatewayService.java @@ -183,7 +183,7 @@ public class GatewayService extends QueueConsumer implements InitializingBean { } OutboundRequest outboundRequest = OutboundRequestBuilder.builder() - .section(Section.MKR.getKey()) + .section(assetOperation.getAmount() != null ? Section.MKR.getKey() : Section.FOND.getKey()) .type(OutboundRequestType.ASSET_OPERATION.getKey()) .content(content).build();