From dcc0a3c7ce8ec0716eb2caaec76cfeab96dba334 Mon Sep 17 00:00:00 2001 From: etreschenkov Date: Fri, 31 May 2024 13:45:06 +0300 Subject: [PATCH] http://jira.mfd.msk:8088/browse/CLS-690 --- .../TradingClearingRegistryService.java | 54 ++++++++++++++----- 1 file changed, 41 insertions(+), 13 deletions(-) diff --git a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/TradingClearingRegistryService.java b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/TradingClearingRegistryService.java index 6ca292696..2d7e01aeb 100644 --- a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/TradingClearingRegistryService.java +++ b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/TradingClearingRegistryService.java @@ -6,6 +6,7 @@ import java.util.Collection; import java.util.List; import java.util.Map; import java.util.Objects; +import java.util.Optional; import java.util.function.Function; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.producer.Producer; @@ -40,6 +41,7 @@ import ru.spcex.clearing.platform.messaging.domain.cud.gateway.Tkr; import ru.spcex.clearing.platform.messaging.domain.cud.registry.TradingClearingRegistryNewRequest; import ru.spcex.clearing.platform.messaging.domain.cud.registry.TradingClearingRegistryUpdateRequest; import ru.spcex.clearing.platform.messaging.domain.cud.reports.NotificationRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest; import ru.spcex.clearing.platform.messaging.serialization.LogFormatter; import ru.spcex.clearing.platform.messaging.service.QueueConsumer; import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate; @@ -49,6 +51,7 @@ import ru.spcex.clearing.util.services.RequestHelper; import ru.spcex.clearing.validation.common.ValidationHelper; import ru.spcex.platform.enumeration.CompanySymbol; import ru.spcex.platform.enumeration.ServiceStatus; +import static ru.spcex.platform.enumeration.Task.makeFiles_MTCR; import ru.spcex.platform.enumeration.TradingClearingRegistryPurpose; import ru.spcex.platform.enumeration.TradingClearingRegistryType; import ru.spcex.platform.imdg.api.Imdg; @@ -142,6 +145,9 @@ public class TradingClearingRegistryService extends QueueConsumer implements Ini callback(TkrAccountsGatewayRequest.class) .setFunction(this::tradingClearingRegistryCheck) .forDestination(Consts.DESTINATION_TRADING_CLEARING_REGISTRY_CHECK, callbacks::put); + callback(LauncherCommandRequest.class) + .setConsumer(this::sendAllTkrToGateway) + .forDestination(makeFiles_MTCR.topic(), callbacks::put); init(); } @@ -229,8 +235,10 @@ public class TradingClearingRegistryService extends QueueConsumer implements Ini targetExistTradingClearingRegistry.setUpdated(Instant.now()); tradingClearingRegistryImdg.update(targetExistTradingClearingRegistry); log.debug("successfully processed, id {}. Updated exist TradingClearingRegistry.id={}", id, targetExistTradingClearingRegistry.getId()); - SendTkrRequest sendTkrRequest = createRequestToGateway(tradingClearingRegistry); - if (!sendTkrRequest.getTkrs().isEmpty()) { + Optional tkr = createRequestToGateway(tradingClearingRegistry); + if (tkr.isPresent()) { + SendTkrRequest sendTkrRequest = new SendTkrRequest(); + sendTkrRequest.getTkrs().add(tkr.get()); kafkaSender.sendRequestToQueue(Consts.ACCOUNTS_TO_GATEWAY, sendTkrRequest); } return null; @@ -367,8 +375,10 @@ public class TradingClearingRegistryService extends QueueConsumer implements Ini tradingClearingRegistryImdg.insert(tradingClearingRegistry); log.info("New TCR.id={} was created.", tradingClearingRegistry.getId()); - SendTkrRequest sendTkrRequest = createRequestToGateway(tradingClearingRegistry); - if (!sendTkrRequest.getTkrs().isEmpty()) { + Optional tkr = createRequestToGateway(tradingClearingRegistry); + if (tkr.isPresent()) { + SendTkrRequest sendTkrRequest = new SendTkrRequest(); + sendTkrRequest.getTkrs().add(tkr.get()); kafkaSender.sendRequestToQueue(Consts.ACCOUNTS_TO_GATEWAY, sendTkrRequest); } @@ -462,8 +472,10 @@ public class TradingClearingRegistryService extends QueueConsumer implements Ini tradingClearingRegistry.setUpdated(now); tradingClearingRegistry.setStatus(req.getStatus()); tradingClearingRegistryImdg.update(tradingClearingRegistry); - SendTkrRequest sendTkrRequest = createRequestToGateway(tradingClearingRegistry); - if (!sendTkrRequest.getTkrs().isEmpty()) { + Optional tkr = createRequestToGateway(tradingClearingRegistry); + if (tkr.isPresent()) { + SendTkrRequest sendTkrRequest = new SendTkrRequest(); + sendTkrRequest.getTkrs().add(tkr.get()); kafkaSender.sendRequestToQueue(Consts.ACCOUNTS_TO_GATEWAY, sendTkrRequest); } log.debug("Update TCR.id={}.", tradingClearingRegistry.getId()); @@ -503,18 +515,31 @@ public class TradingClearingRegistryService extends QueueConsumer implements Ini TradingClearingRegistry tradingClearingRegistry = tradingClearingRegistryImdg.getSingleObjectByFieldValues( Map.of("code", accountForCheck.getTkrCode()) ); - SendTkrRequest sendTkrRequest = createRequestToGateway(tradingClearingRegistry); - sendTkrRequest.setRequestId(req.getRequestId()); - if (!sendTkrRequest.getTkrs().isEmpty()) { + Optional tkr = createRequestToGateway(tradingClearingRegistry); + if (tkr.isPresent()) { + SendTkrRequest sendTkrRequest = new SendTkrRequest(); + sendTkrRequest.getTkrs().add(tkr.get()); kafkaSender.sendRequestToQueue(Consts.ACCOUNTS_TO_GATEWAY, sendTkrRequest); } } return null; } - private SendTkrRequest createRequestToGateway(TradingClearingRegistry tradingClearingRegistry) { + public RequestInfoUpdate sendAllTkrToGateway(BaseRequest gatewayRequest) { + Collection tkrs = tradingClearingRegistryImdg.getAllValues(); SendTkrRequest sendTkrRequest = new SendTkrRequest(); + for (TradingClearingRegistry tkr : tkrs) { + log.debug("Create TKR request with id: {}", tkr.getId()); + Optional tkrRequest = createRequestToGateway(tkr); + tkrRequest.ifPresent(value -> sendTkrRequest.getTkrs().add(value)); + } + kafkaSender.sendRequestToQueue(Consts.ACCOUNTS_TO_GATEWAY, sendTkrRequest); + return null; + } + + private Optional createRequestToGateway(TradingClearingRegistry tradingClearingRegistry) { if (tradingClearingRegistry != null) { + Tkr tkr = new Tkr(); log.debug("Found TKR with id: {} and code: {}", tradingClearingRegistry.getId(), tradingClearingRegistry.getCode()); Company company = companyImdg.getSingleObjectByID(tradingClearingRegistry.getCompanyId()); CompanySymbols companySymbols = companySymbolsImdg.getSingleObjectByFieldValues( @@ -523,10 +548,13 @@ public class TradingClearingRegistryService extends QueueConsumer implements Ini "companyId", company.getId() ) ); + if (companySymbols == null){ + log.warn("CompanySymbols is null, searched by company.id: {}, skip this tkr", company.getId()); + return Optional.empty(); + } ClientCode clientCode = clientCodeImdg.getFirstObjectByFieldValues( Map.of("tradingClearingRegistryId", tradingClearingRegistry.getId()) ); - Tkr tkr = new Tkr(); tkr.setCompanyId(companySymbols.getCompanySymbolValue()); tkr.setTradingCode(Long.valueOf(company.getTradingCode())); tkr.setTkrCode(tradingClearingRegistry.getCode()); @@ -575,9 +603,9 @@ public class TradingClearingRegistryService extends QueueConsumer implements Ini tkr.getMoneyAccounts().add(moneyAccountMsg); } } - sendTkrRequest.getTkrs().add(tkr); + return Optional.of(tkr); } - return sendTkrRequest; + return Optional.empty(); } /**