diff --git a/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/services/AbstractExporterService.java b/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/services/AbstractExporterService.java index aa3a46ede..3547187f0 100644 --- a/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/services/AbstractExporterService.java +++ b/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/services/AbstractExporterService.java @@ -7,10 +7,14 @@ import ru.clearing.classes.statics.data.registry.Registry; import ru.clearing.classes.statics.data.registry.TradingClearingRegistry; import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.lim.exporter.config.SFTPConfig; +import ru.spcex.clearing.platform.messaging.domain.Consts; import ru.spcex.clearing.platform.messaging.domain.cud.utilities.LimExportedRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.utilities.NotificationNewRequest; import ru.spcex.clearing.platform.messaging.serialization.LogFormatter; import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate; import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; +import ru.spcex.platform.enumeration.ObjectType; +import ru.spcex.platform.enumeration.Priority; import ru.spcex.platform.enumeration.Section; import ru.spcex.platform.enumeration.ServiceStatus; import ru.spcex.platform.imdg.api.Imdg; @@ -65,6 +69,7 @@ public abstract class AbstractExporterService { gateway.sendToSftp(limFile); } catch (IOException e) { log.error("Failed export {} file", fileName); + sendErrorNotification(fileName); throw new RuntimeException(e); } finally { try { @@ -87,6 +92,15 @@ public abstract class AbstractExporterService { kafkaSender.sendRequestToQueue(LIM_EXPORTED, limExportedRequest); } + void sendErrorNotification(String fileName) { + NotificationNewRequest newRequest = new NotificationNewRequest(); + newRequest.setObjectType(ObjectType.rgst.getKey()); + newRequest.setComment("Failed export file = " + fileName + " for section = " + section().getKey()); + newRequest.setPriority(Priority.HIGH.getKey()); + log.info("Sending notification to kafka: {}", newRequest); + kafkaSender.sendRequestToQueue(Consts.NOTIFICATION_NEW, newRequest); + } + protected String prepareFileName(String target) { StringBuilder result = new StringBuilder(); String dt = dtFormatter.format(LocalDateTime.now()); diff --git a/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/services/CurrencyExporterService.java b/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/services/CurrencyExporterService.java new file mode 100644 index 000000000..0fe1ac4e1 --- /dev/null +++ b/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/services/CurrencyExporterService.java @@ -0,0 +1,97 @@ +package ru.spcex.clearing.lim.exporter.services; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.stereotype.Service; +import ru.clearing.classes.statics.data.registry.Registry; +import ru.spcex.clearing.lim.exporter.config.LimFormat; +import ru.spcex.clearing.lim.exporter.config.SFTPConfig; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; +import ru.spcex.platform.enumeration.RegistryStatus; +import ru.spcex.platform.enumeration.Section; +import ru.spcex.platform.imdg.api.ImdgProvider; +import ru.spcex.platform.imdg.api.predicate.ImdgPredicate; +import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder; +import ru.spcex.platform.imdg.api.predicate.specific.RegistryCodeSqlBuilder; + +import java.math.BigDecimal; +import java.util.ArrayList; +import java.util.Collection; +import java.util.List; +import java.util.function.Function; + +import static ru.spcex.platform.enumeration.RegistryTradingParams.*; + +@Service +public class CurrencyExporterService extends AbstractExporterService { + private final Logger log = LoggerFactory.getLogger(getClass()); + + public CurrencyExporterService(SFTPConfig.LimGateway gateway, + KafkaSender kafkaSender, + ImdgProvider imdgProvider) { + super(gateway, kafkaSender, imdgProvider); + } + + @Override + public String getTargetFileName() { + return prepareFileName("money"); + } + + @Override + public Section section() { + return Section.CURR; + } + + @Override + public Collection getLimFileRows() { + log.debug("Started loading and formation of currency(money) file lines"); + List limFileRows = new ArrayList<>(); + Collection registriesA = registryImdg.getCollectionObjectsBySQL(RegistryCodeSqlBuilder.getInstance(AM_F).build()); + log.debug("Selected {} registriesA.", registriesA.size()); + + ImdgPredicateBuilder pb = registryImdg.predicateBuilder(); + Function prdctDM_T = registry -> pb.and( + pb.sql(RegistryCodeSqlBuilder.getInstance(DM_T).build()), + pb.equals("registryStatus", RegistryStatus.PROC.getKey()), + pb.equals("tradingClearingRegistry", registry.getTradingClearingRegistry()), + pb.equals("securityId", registry.getSecurityId()) + ); + for (Registry registry : registriesA) { + if (checkNotBlocked(registry)) { + BigDecimal sumPositive = BigDecimal.ZERO; + if (registry.getTradingClearingRegistry() != null && registry.getSecurityId() != null) { + Collection regsDMT = registryImdg.getCollectionObjectsByPredicate(prdctDM_T.apply(registry)); + log.trace("Found {} DM_T registries for registry AM_F {}", regsDMT.size(), registry); + sumPositive = sumPositive.add(sum(regsDMT)); + } + limFileRows.add(getRow(registry, sumPositive)); + } + } + log.debug("Successfully completed the formation of rows: {} for export currency(money)", limFileRows.size()); + return limFileRows; + } + + public String getRow(Registry registryA, BigDecimal sumPositive) { + StringBuilder row = new StringBuilder(); + + row.append("MONEY: FIRM_ID = "); + row.append(registryA.getTradingCode()); + + row.append("; TAG = SPVB"); + + row.append("; CURR_CODE = "); + row.append(registryA.getSecuritySymbol()); + + row.append("; CLIENT_CODE = "); + row.append(registryA.getTradingClearingRegistry()); + + row.append("; OPEN_BALANCE = "); + BigDecimal balance = registryA.getBalance() != null ? registryA.getBalance().subtract(sumPositive) : BigDecimal.ZERO; + row.append(LimFormat.toStringD2(balance)); + + row.append("; OPEN_LIMIT = 0.00"); + + row.append("; LIMIT_KIND = 0;"); + return row.toString(); + } +} diff --git a/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/services/LauncherCommandReceiver.java b/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/services/LauncherCommandReceiver.java index 9dea4ab50..aaee769fc 100644 --- a/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/services/LauncherCommandReceiver.java +++ b/clearing-parent/lim-exporter/src/main/java/ru/spcex/clearing/lim/exporter/services/LauncherCommandReceiver.java @@ -15,14 +15,17 @@ public class LauncherCommandReceiver extends QueueConsumer implements Initializi private final Logger log = LoggerFactory.getLogger(getClass()); private final MoneyExporterService moneyExporterService; private final SecurityExporterService securityExporterService; + private final CurrencyExporterService currencyExporterService; public LauncherCommandReceiver(Consumer kafkaQueue, Producer kafkaProducer, MoneyExporterService moneyExporterService, - SecurityExporterService securityExporterService) { + SecurityExporterService securityExporterService, + CurrencyExporterService currencyExporterService) { super(kafkaQueue, kafkaProducer); this.moneyExporterService = moneyExporterService; this.securityExporterService = securityExporterService; + this.currencyExporterService = currencyExporterService; } @Override @@ -33,6 +36,9 @@ public class LauncherCommandReceiver extends QueueConsumer implements Initializi callback(LauncherCommandRequest.class) .setFunction(action -> securityExporterService.process()) .forDestination(Task.unloadingSession_LIMS.topic(), callbacks::put); // LIMS + callback(LauncherCommandRequest.class) + .setFunction(action -> currencyExporterService.process()) + .forDestination(Task.unloadingSession_LIMC.topic(), callbacks::put); // LIMC init(); } } diff --git a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/Task.java b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/Task.java index 9786862b9..d316c0c6e 100644 --- a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/Task.java +++ b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/Task.java @@ -23,6 +23,7 @@ public enum Task implements IEnumKey { createDealsReport_GREF("GREF"), // Формирование финальной отчетности по сделкам unloadingSession_LIMM("LIMM"),//Выгрузка в торговую систему остатков по деньгам unloadingSession_LIMS("LIMS"),//Выгрузка в торговую систему остатков по бумагам + unloadingSession_LIMC("LIMC"),//Выгрузка в торговую систему остатков по валютам liquidationSession_LIQU("LIQU"),//Ликвидационная сессия по обязательтсвам участника startSession_STRM("STRM"),//Начало торговой сессии секции МКР terminationSession_ETRM("ETRM"),//Завершение торговой сессии секции МКР