Compare commits

..

1 commit

Author SHA1 Message Date
ialbert
56c7573fd3 IKafkaRequestStatusProcessor implementation 2023-12-04 19:17:41 +03:00
51 changed files with 476 additions and 746 deletions

View file

@ -70,7 +70,6 @@ public class ClearingAccountService extends QueueConsumer implements Initializin
private final Function<ClearingAccountUpdateRequest, IValidator> clearingAccountUpdateRequestValidator;
private final Imdg<Account> accountImdg;
private final Imdg<ClearingAccount> clearingAccountImdg;
private final Imdg<Company> companyImdg;
private final Imdg<Relation> relationImdg;
private final Imdg<ClearingMemberCategory> clearingMemberCategoryImdg;
@ -104,7 +103,6 @@ public class ClearingAccountService extends QueueConsumer implements Initializin
this.clearingAccountUpdateRequestValidator = clearingAccountUpdateRequestValidator;
this.accountImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Account, Account.class);
this.clearingAccountImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ClearingAccount, ClearingAccount.class);
this.companyImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Company, Company.class);
this.relationImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Relation, Relation.class);
this.clearingMemberCategoryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ClearingMemberCategory, ClearingMemberCategory.class);
@ -364,16 +362,11 @@ public class ClearingAccountService extends QueueConsumer implements Initializin
);
Account account = sDf52.getAccount() == null ? null : accountImdg.getFirstObjectByFieldValues(accountQuery);
if (account == null) {
if (SDFProcessService.SDF52_STATUS_3Open.equals(sDf52.getStatus())) {
log.debug("By generationId={} s_df52[{}].status={}, but account not found (query: {}). COntinuse with result OK for status 3",
groupId, sDf52.getId(), sDf52.getStatus(), accountQuery);
} else {
String msg = messageResolver.resolve(new EnumMessage(AccountError.AccountNotFound, sDf52.getAccount()));
log.warn("By generationId={} s_df52[{}] (query: {}) error: {}", groupId, sDf52.getId(), accountQuery, msg);
toProcessSDF53.add(new MutableTriple<>(sDf52, null,
makeSdfErrorText(AccountError.AccountNotFound, SDFProcessService.SDF_STATUS_ERROR_COMPANY_NOT_FOUND)));
continue;
}
String msg = messageResolver.resolve(new EnumMessage(AccountError.AccountNotFound, sDf52.getAccount()));
log.warn("By generationId={} s_df52[{}] (query: {}) error: {}", groupId, sDf52.getId(), accountQuery, msg);
toProcessSDF53.add(new MutableTriple<>(sDf52, null,
makeSdfErrorText(AccountError.AccountNotFound, SDFProcessService.SDF_STATUS_ERROR_COMPANY_NOT_FOUND)));
continue;
}
toProcessSDF53.add(new MutableTriple<>(sDf52, account, SDFProcessService.SDF_STATUS_OK));
toUpdate.add(new Pair<>(sDf52, account));

View file

@ -38,11 +38,6 @@ public class SDFProcessService {
public static final String SDF_STATUS_ERROR_NO_TRADE = "4"; // (Ошибка. Торги не идут)
public static final String SDF_STATUS_ERROR = "9"; // (Другие ошибки, выявленные в КС) - если получены другие ошибки (т.е. account НЕ обновлена по sDf52).
protected static final Long SDF52_STATUS_0Blocked = 0L;
protected static final Long SDF52_STATUS_1Unblocked = 1L;
protected static final Long SDF52_STATUS_2Closed = 2L;
protected static final Long SDF52_STATUS_3Open = 3L;
final protected Logger log = LoggerFactory.getLogger(getClass());
protected final Producer<String, Object> kafkaResponseQueue;
@ -134,12 +129,12 @@ public class SDFProcessService {
* @return AccountStatus или null
*/
public AccountStatus parseSdf52Status(Long status) {
if (SDF52_STATUS_1Unblocked.equals(status) || SDF52_STATUS_3Open.equals(status)) {
if (Long.valueOf(1L).equals(status) || Long.valueOf(3L).equals(status)) {
return AccountStatus.ACTIVE;
} else if (SDF52_STATUS_0Blocked.equals(status)) {
} else if (Long.valueOf(0L).equals(status)) {
return AccountStatus.BLOCKED;
}
if (SDF52_STATUS_2Closed.equals(status)) {
if (Long.valueOf(2L).equals(status)) {
return AccountStatus.CLOSE;
}
return null;

View file

@ -1,6 +1,6 @@
{
"version": "3.9.0.71",
"version": "3.9.0.70",
"enums": {
@ -6182,11 +6182,11 @@
"fields": [
{"code": "id",
"type": 1,"name": "Идентификатор записи","shortname": "ID","searchable": true,"sortable": true,"ignore": true
"type": 1,"name": "Идентификатор записи","shortname": "ID","searchable": true,"sortable": true
}
,
{"code": "createdAt",
"field": "created","type": 4,"webtype": "5","dbname": "Дата-время создания записи","name": "Время создания записи","shortname": "Создано","searchable": true,"sortable": true
"field": "created","type": 4,"webtype": "5","dbname": "Дата-время создания записи","name": "Время создания записи","shortname": "Создано","searchable": true,"sortable": true,"ignore": true
}
,
{"code": "updatedAt",
@ -6198,7 +6198,7 @@
}
,
{"code": "sessionId",
"type": 1,"dbname": "Идентификатор клиринговой сессии","name": "Клиринговая сессия","shortname": "Сессия","visible": true,"searchable": true,"sortable": true,"link": "session","ignore": true
"type": 1,"dbname": "Идентификатор клиринговой сессии","name": "Клиринговая сессия","shortname": "Сессия","visible": true,"searchable": true,"sortable": true,"link": "session"
}
,
{"code": "depoCode",
@ -6238,7 +6238,7 @@
}
,
{"code": "infoAccount",
"type": 2,"length": 50,"name": "Номер счета внутреннего учета СПВБ","shortname": "Номер счета внутреннего учета СПВБ","searchable": true,"sortable": true,"visible": true,"ignore": true
"type": 2,"length": 50,"name": "Номер счета внутреннего учета СПВБ","shortname": "Номер счета внутреннего учета СПВБ","searchable": true,"sortable": true,"visible": true
}
,
{"code": "remainderSum",
@ -6246,23 +6246,23 @@
}
,
{"code": "blockedSum",
"type": 10,"name": "Сумма блокированных денежных средств","shortname": "Блокированные","searchable": true,"sortable": true,"ignore": true
"type": 10,"name": "Сумма блокированных денежных средств","shortname": "Блокированные","searchable": true,"sortable": true
}
,
{"code": "unblockedSum",
"type": 10,"name": "Сумма свободных денежных средств","shortname": "Свободные","searchable": true,"sortable": true,"ignore": true
"type": 10,"name": "Сумма свободных денежных средств","shortname": "Свободные","searchable": true,"sortable": true
}
,
{"code": "inn",
"type": 2,"length": 255,"name": "Идентификационный номер налогоплательщика (ИНН)","shortname": "ИНН","searchable": true,"sortable": true,"visible": true,"ignore": true
"type": 2,"length": 255,"name": "Идентификационный номер налогоплательщика (ИНН)","shortname": "ИНН","searchable": true,"sortable": true,"visible": true
}
,
{"code": "sessionId",
"type": 1,"dbname": "Идентификатор клиринговой сессии","name": "Клиринговая сессия","shortname": "Сессия","visible": true,"searchable": true,"sortable": true,"link": "session","ignore": true
"type": 1,"dbname": "Идентификатор клиринговой сессии","name": "Клиринговая сессия","shortname": "Сессия","visible": true,"searchable": true,"sortable": true,"link": "session"
}
,
{"code": "companyFullName",
"type": 2,"length": 255,"name": "Полное наименование компании","shortname": "Полное наименование компании","searchable": true,"sortable": true,"visible": true,"ignore": true
"type": 2,"length": 255,"name": "Полное наименование компании","shortname": "Полное наименование компании","searchable": true,"sortable": true,"visible": true
}
,
{"code": "companyId",
@ -6270,7 +6270,7 @@
}
,
{"code": "id",
"type": 1,"name": "Идентификатор записи","shortname": "ID","searchable": true,"sortable": true,"ignore": true
"type": 1,"name": "Идентификатор записи","shortname": "ID","searchable": true,"sortable": true
}
,
{"code": "createdAt",
@ -6278,7 +6278,7 @@
}
,
{"code": "updatedAt",
"field": "updated","type": 4,"webtype": "5","dbname": "Дата-время изменения записи","name": "Время изменения записи","shortname": "Изменено","searchable": true,"sortable": true,"ignore": true
"field": "updated","type": 4,"webtype": "5","dbname": "Дата-время изменения записи","name": "Время изменения записи","shortname": "Изменено","searchable": true,"sortable": true
}
]

View file

@ -1,6 +1,6 @@
<?xml version="1.0" encoding="UTF-8"?>
<!--?xml-stylesheet type="text/xsl" href="\..\corp-reports\src\data\meta\meta.server.xslt"?-->
<meta version="3.9.0.71">
<meta version="3.9.0.70">
<!-- _xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" _xsi:noNamespaceSchemaLocation="file:///E:/d/projects/meta/from/meta.xsd" -->
<!--Здесь словари-->
<enums>
@ -1427,11 +1427,11 @@
<sessionId type="1" dbname="Идентификатор клиринговой сессии" name="Клиринговая сессия" shortname="Сессия" visible="true" searchable="true" sortable="true" link="session"/>
</executionFond>
<depoBalanceRegister name="Реестр остатков ценных бумаг" destination="depo-balance-registers" historyDestination="history" class="ru.clearing.classes.statics.data.register.DepoBalanceRegister" table="balance_depo_register">
<id type="1" name="Идентификатор записи" shortname="ID" searchable="true" sortable="true" ignore="true"/>
<createdAt field="created" type="4" webtype="5" dbname="Дата-время создания записи" name="Время создания записи" shortname="Создано" searchable="true" sortable="true"/>
<id type="1" name="Идентификатор записи" shortname="ID" searchable="true" sortable="true"/>
<createdAt field="created" type="4" webtype="5" dbname="Дата-время создания записи" name="Время создания записи" shortname="Создано" searchable="true" sortable="true" ignore="true"/>
<updatedAt field="updated" type="4" webtype="5" dbname="Дата-время изменения записи" name="Время изменения записи" shortname="Изменено" searchable="true" sortable="true" ignore="true"/>
<companyId type="1" dbname="Идентификатор компании" name="Наименование компании" shortname="Компания" visible="true" searchable="true" sortable="true" link="company" linkCode="shortName"/>
<sessionId type="1" dbname="Идентификатор клиринговой сессии" name="Клиринговая сессия" shortname="Сессия" visible="true" searchable="true" sortable="true" link="session" ignore="true"/>
<sessionId type="1" dbname="Идентификатор клиринговой сессии" name="Клиринговая сессия" shortname="Сессия" visible="true" searchable="true" sortable="true" link="session"/>
<depoCode type="2" length="50" name="Код раздела субсчета/счета депо" shortname="Код счета депо" searchable="true" sortable="true"/>
<quantity type="11" name="Количество" shortname="Количество" searchable="true" sortable="true"/>
<securitySymbol type="2" length="255" name="Код ценной бумаги" shortname="Ценная бумага" searchable="true" sortable="true"/>
@ -1439,17 +1439,17 @@
<moneyBalanceRegister name="Реестр остатков денежных средств" destination="money-balance-registers" historyDestination="history" class="ru.clearing.classes.statics.data.register.MoneyBalanceRegister" table="money_balance_register">
<setHouseName type="2" length="255" name="Наименование РО" shortname="Наименование РО" searchable="true" sortable="true" visible="true"/>
<account type="2" length="50" name="Номер торгового/клирингового счета" shortname="Номер торгового/клирингового счета" searchable="true" sortable="true" visible="true"/>
<infoAccount type="2" length="50" name="Номер счета внутреннего учета СПВБ" shortname="Номер счета внутреннего учета СПВБ" searchable="true" sortable="true" visible="true" ignore="true"/>
<infoAccount type="2" length="50" name="Номер счета внутреннего учета СПВБ" shortname="Номер счета внутреннего учета СПВБ" searchable="true" sortable="true" visible="true"/>
<remainderSum type="10" name="Остаток денежных средств" shortname="Остаток" searchable="true" sortable="true"/>
<blockedSum type="10" name="Сумма блокированных денежных средств" shortname="Блокированные" searchable="true" sortable="true" ignore="true"/>
<unblockedSum type="10" name="Сумма свободных денежных средств" shortname="Свободные" searchable="true" sortable="true" ignore="true"/>
<inn type="2" length="255" name="Идентификационный номер налогоплательщика (ИНН)" shortname="ИНН" searchable="true" sortable="true" visible="true" ignore="true"/>
<sessionId type="1" dbname="Идентификатор клиринговой сессии" name="Клиринговая сессия" shortname="Сессия" visible="true" searchable="true" sortable="true" link="session" ignore="true"/>
<companyFullName type="2" length="255" name="Полное наименование компании" shortname="Полное наименование компании" searchable="true" sortable="true" visible="true" ignore="true"/>
<blockedSum type="10" name="Сумма блокированных денежных средств" shortname="Блокированные" searchable="true" sortable="true"/>
<unblockedSum type="10" name="Сумма свободных денежных средств" shortname="Свободные" searchable="true" sortable="true"/>
<inn type="2" length="255" name="Идентификационный номер налогоплательщика (ИНН)" shortname="ИНН" searchable="true" sortable="true" visible="true"/>
<sessionId type="1" dbname="Идентификатор клиринговой сессии" name="Клиринговая сессия" shortname="Сессия" visible="true" searchable="true" sortable="true" link="session"/>
<companyFullName type="2" length="255" name="Полное наименование компании" shortname="Полное наименование компании" searchable="true" sortable="true" visible="true"/>
<companyId type="1" dbname="Идентификатор компании" name="Наименование компании" shortname="Компания" searchable="true" sortable="true" visible="true" link="company" linkCode="shortName"/>
<id type="1" name="Идентификатор записи" shortname="ID" searchable="true" sortable="true" ignore="true"/>
<id type="1" name="Идентификатор записи" shortname="ID" searchable="true" sortable="true"/>
<createdAt field="created" type="4" webtype="5" dbname="Дата-время создания записи" name="Время создания записи" shortname="Создано" searchable="true" sortable="true"/>
<updatedAt field="updated" type="4" webtype="5" dbname="Дата-время изменения записи" name="Время изменения записи" shortname="Изменено" searchable="true" sortable="true" ignore="true"/>
<updatedAt field="updated" type="4" webtype="5" dbname="Дата-время изменения записи" name="Время изменения записи" shortname="Изменено" searchable="true" sortable="true"/>
</moneyBalanceRegister>
<admittedLiabilitiesRegister name="Реестр обязательств, допущенных к клирингу" destination="admitted-liabilities-registers" historyDestination="history" class="ru.clearing.classes.statics.data.register.AdmittedLiabilitiesRegister" table="admitted_liabilities_register">
<companyFullName type="2" length="255" name="Полное наименование компании" shortname="Полное наименование компании" searchable="true" sortable="true" visible="true"/>

View file

@ -1,6 +1,6 @@
{
"version": "3.9.0.71",
"version": "3.9.0.70",
"enums": {
@ -6182,11 +6182,11 @@
"fields": [
{"code": "id",
"type": 1,"name": "Идентификатор записи","shortname": "ID","searchable": true,"sortable": true,"ignore": true
"type": 1,"name": "Идентификатор записи","shortname": "ID","searchable": true,"sortable": true
}
,
{"code": "createdAt",
"field": "created","type": 4,"webtype": "5","dbname": "Дата-время создания записи","name": "Время создания записи","shortname": "Создано","searchable": true,"sortable": true
"field": "created","type": 4,"webtype": "5","dbname": "Дата-время создания записи","name": "Время создания записи","shortname": "Создано","searchable": true,"sortable": true,"ignore": true
}
,
{"code": "updatedAt",
@ -6198,7 +6198,7 @@
}
,
{"code": "sessionId",
"type": 1,"dbname": "Идентификатор клиринговой сессии","name": "Клиринговая сессия","shortname": "Сессия","visible": true,"searchable": true,"sortable": true,"link": "session","ignore": true
"type": 1,"dbname": "Идентификатор клиринговой сессии","name": "Клиринговая сессия","shortname": "Сессия","visible": true,"searchable": true,"sortable": true,"link": "session"
}
,
{"code": "depoCode",
@ -6238,7 +6238,7 @@
}
,
{"code": "infoAccount",
"type": 2,"length": 50,"name": "Номер счета внутреннего учета СПВБ","shortname": "Номер счета внутреннего учета СПВБ","searchable": true,"sortable": true,"visible": true,"ignore": true
"type": 2,"length": 50,"name": "Номер счета внутреннего учета СПВБ","shortname": "Номер счета внутреннего учета СПВБ","searchable": true,"sortable": true,"visible": true
}
,
{"code": "remainderSum",
@ -6246,23 +6246,23 @@
}
,
{"code": "blockedSum",
"type": 10,"name": "Сумма блокированных денежных средств","shortname": "Блокированные","searchable": true,"sortable": true,"ignore": true
"type": 10,"name": "Сумма блокированных денежных средств","shortname": "Блокированные","searchable": true,"sortable": true
}
,
{"code": "unblockedSum",
"type": 10,"name": "Сумма свободных денежных средств","shortname": "Свободные","searchable": true,"sortable": true,"ignore": true
"type": 10,"name": "Сумма свободных денежных средств","shortname": "Свободные","searchable": true,"sortable": true
}
,
{"code": "inn",
"type": 2,"length": 255,"name": "Идентификационный номер налогоплательщика (ИНН)","shortname": "ИНН","searchable": true,"sortable": true,"visible": true,"ignore": true
"type": 2,"length": 255,"name": "Идентификационный номер налогоплательщика (ИНН)","shortname": "ИНН","searchable": true,"sortable": true,"visible": true
}
,
{"code": "sessionId",
"type": 1,"dbname": "Идентификатор клиринговой сессии","name": "Клиринговая сессия","shortname": "Сессия","visible": true,"searchable": true,"sortable": true,"link": "session","ignore": true
"type": 1,"dbname": "Идентификатор клиринговой сессии","name": "Клиринговая сессия","shortname": "Сессия","visible": true,"searchable": true,"sortable": true,"link": "session"
}
,
{"code": "companyFullName",
"type": 2,"length": 255,"name": "Полное наименование компании","shortname": "Полное наименование компании","searchable": true,"sortable": true,"visible": true,"ignore": true
"type": 2,"length": 255,"name": "Полное наименование компании","shortname": "Полное наименование компании","searchable": true,"sortable": true,"visible": true
}
,
{"code": "companyId",
@ -6270,7 +6270,7 @@
}
,
{"code": "id",
"type": 1,"name": "Идентификатор записи","shortname": "ID","searchable": true,"sortable": true,"ignore": true
"type": 1,"name": "Идентификатор записи","shortname": "ID","searchable": true,"sortable": true
}
,
{"code": "createdAt",
@ -6278,7 +6278,7 @@
}
,
{"code": "updatedAt",
"field": "updated","type": 4,"webtype": "5","dbname": "Дата-время изменения записи","name": "Время изменения записи","shortname": "Изменено","searchable": true,"sortable": true,"ignore": true
"field": "updated","type": 4,"webtype": "5","dbname": "Дата-время изменения записи","name": "Время изменения записи","shortname": "Изменено","searchable": true,"sortable": true
}
]

View file

@ -2,8 +2,6 @@ package ru.spcex.clearing.config;
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.producer.Producer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.config.ConfigurableBeanFactory;
import org.springframework.context.annotation.Bean;
@ -28,31 +26,19 @@ import java.util.function.Supplier;
@Configuration
public class KafkaConfig {
private final Logger log = LoggerFactory.getLogger(getClass());
private final KafkaConsumerSettings kafkaSettings;
public KafkaConfig(ClearingServiceSettings settings) {
this.kafkaSettings = settings.getKafkaConsumer();
Integer gtwTimeout = settings.getSessionStage().getInspectionGatewayTimeout();
if (gtwTimeout != null) {
int maxPollIntervalMs = gtwTimeout * 1000 + 1000;
log.debug("setting max.poll.interval.ms {}", maxPollIntervalMs);
this.kafkaSettings.setManuallyMaxPollIntervalMs(maxPollIntervalMs);
}
}
@Autowired
@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE)
@Bean("kafkaConsumer")
public Consumer<String, Object> createConsumer() {
return KafkaConsumerFactory.consumer(kafkaSettings);
public Consumer<String, Object> createConsumer(ClearingServiceSettings settings) {
return KafkaConsumerFactory.consumer(settings.getKafkaConsumer());
}
@Autowired
@Bean("kafkaConsumerGateway")
public Supplier<Consumer<String, Object>> getwaySessionConsumer() {
public Supplier<Consumer<String, Object>> getwaySessionConsumer(ClearingServiceSettings settings) {
AtomicInteger groupId = new AtomicInteger(1);
return () -> {
KafkaConsumerSettings consumer = kafkaSettings;
KafkaConsumerSettings consumer = settings.getKafkaConsumer();
KafkaConsumerSettings gatewayConsumer = new KafkaConsumerSettings();
gatewayConsumer.setBootstrapServers(consumer.getBootstrapServers());
gatewayConsumer.setSessionTimeoutMs(consumer.getSessionTimeoutMs());

View file

@ -368,7 +368,6 @@ public class ValidationConfig {
new PresentById(IMDGDistributedNames.Map_Registry, ClearingError.RecordNotFound, true),
StatusExtractValidationRule.RegistryCodeCheck,
StatusExtractValidationRule.ContractCheck,
StatusExtractValidationRule.RequestStatusValid,
StatusExtractValidationRule.RegistryStatusValid
);
};

View file

@ -48,11 +48,16 @@ import ru.spcex.platform.utils.validation.IValidator;
import java.math.BigDecimal;
import java.time.Instant;
import java.time.LocalDate;
import java.util.*;
import java.util.Collection;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.function.Function;
import java.util.function.Supplier;
import java.util.stream.Collectors;
import java.util.stream.Stream;
import static ru.spcex.clearing.session.stage.impl.GatewayRequester.mapError;
import static ru.spcex.platform.utils.number.BigDecimalUtil.safeBD;
@Service
@ -285,11 +290,23 @@ public class RegistryService {
return reqHelp.error(req.getId(), err.get());
}
Registry rgs = validator.getStored(Stored.PresentById);
Supplier<AssetOperationRequest> gtwBuilder = () -> GatewayRequestCreator.from(rgs, InOutDirection.in);
Optional<Boolean> gatewayOk;
if (trdTime.isTradingTime()) {
Long gtwReqId = kafkaSender.sendRequestToQueue(Consts.ASSET_OPERATION, gtwReq(rgs));
log.debug("send request to {} id={}", Consts.ASSET_OPERATION, gtwReqId);
gatewayOk = gateway.gatewayRequestAndWait(gtwBuilder);
} else {
log.debug("RegistryChangeStatusExtractRequest rgs.id={} not sending gateway request", rgs.getId());
gatewayOk = Optional.of(true);
}
if (gatewayOk.isEmpty() || !gatewayOk.get()) {
String gtwErr = msgResolver.resolve(mapError(gatewayOk));
log.error("{}.id={} {}",
rgs.getRegistryCode(),
rgs.getId(),
gtwErr);
notification.sendNotification(ObjectType.rgst,
"Отметка о получении выписки: %s".formatted(gtwErr),
Priority.HIGH);
}
rgs.setRegistryStatus(payload.getRegistryStatus());
log.debug("changing registry.id={} status to {}", rgs.getId(), payload.getRegistryStatus());
@ -304,13 +321,6 @@ public class RegistryService {
return reqHelp.success(req.getId());
}
private AssetOperationListRequest gtwReq(Registry rgs) {
AssetOperationRequest item = GatewayRequestCreator.from(rgs, InOutDirection.in);
AssetOperationListRequest gtwReq = new AssetOperationListRequest();
gtwReq.setAssetOperationRequests(Collections.singletonList(item));
return gtwReq;
}
public RequestInfoUpdate identificationFunds(BaseRequest<IdentificationFundsRequest> req) {
if (!rights.userHasRole(req.getUserId(), UserRole.Admin)) {
return reqHelp.error(req.getId(), ClearingError.UserVerifyDenial);

View file

@ -236,7 +236,7 @@ public class PaymentInstructionBuilderV2 {
boolean hasCounterPartyInitiator = false;
for (Registry claim : claims) {
Long counterPartyId = claim.getCounterPartyId();
if (categoryIOrVInitiator(counterPartyId) && companyRoleODEPPresent(counterPartyId)) {
if (categoryVInitiator(counterPartyId) && companyRoleODEPPresent(counterPartyId)) {
hasCounterPartyInitiator = true;
break;
}
@ -244,11 +244,11 @@ public class PaymentInstructionBuilderV2 {
return hasCounterPartyInitiator;
}
public boolean categoryIOrVInitiator(Long companyId) {
public boolean categoryVInitiator(Long companyId) {
ClearingMemberCategory category = clearingCategoryImdg.getFirstObjectByFieldValues(
Map.of("companyId", companyId));
ClearingCategory ctg = IEnumKey.getEnumByKey(ClearingCategory.class, category.getClearingMemberCategory());
return ctg != null && (ctg.equals(ClearingCategory.V) || ctg.equals(ClearingCategory.I)) ;
return ctg != null && ctg.equals(ClearingCategory.V);
}
public boolean companyRoleODEPPresent(Long companyId) {

View file

@ -133,7 +133,6 @@ public class RegistryBuilder {
rgs.setCredit(BigDecimal.ZERO);
rgs.setDiffBalance(BigDecimal.ZERO);
rgs.setCheckBalance(BigDecimal.ZERO);
rgs.setOpenBalance(BigDecimal.ZERO);
return rgs;
}

View file

@ -133,7 +133,6 @@ public class RegistrySecurityBuilder {
rgs.setCredit(BigDecimal.ZERO);
rgs.setDiffBalance(BigDecimal.ZERO);
rgs.setCheckBalance(BigDecimal.ZERO);
rgs.setOpenBalance(BigDecimal.ZERO);
return rgs;
}

View file

@ -346,7 +346,6 @@ public class Sdf08Executor extends AbstractExecutor<SDf08> {
rgs.setCredit(BigDecimal.ZERO);
rgs.setDiffBalance(BigDecimal.ZERO);
rgs.setCheckBalance(BigDecimal.ZERO);
rgs.setOpenBalance(BigDecimal.ZERO);
rgs.setBalanceDimension(BalanceDimension.MONY.getKey());
rgs.setTradingDate(statement.getSettlementDate());
rgs.setClearingDate(LocalDate.now());

View file

@ -169,9 +169,6 @@ public class Sdf57Executor extends AbstractExecutor<SDf57> {
result.setChildGenerationId(generationIdForGroup);
log.info("SDF57 execution: sdf57 number={}, groupId={}", sdf.size(), sdf.stream().findFirst().map(SDf57::getGenerationId).orElse(null));
AtomicReference<Boolean> sessionIsNeededFlag = new AtomicReference<>(false);
boolean sessionIsPresent = sessionImdg.getFirstObjectByFieldValues(Map.of(
"workflowStatus", SessionStatus.ACTV.getKey()
)) != null;
for (SDf57 sdf57 : sdf) {
IValidator validator = sDf57Validator.apply(sdf57);
Optional<EnumMessage> error = validator.tillFirstError();
@ -366,7 +363,8 @@ public class Sdf57Executor extends AbstractExecutor<SDf57> {
if (TextUtil.isEmpty(asts.a__t().getTradingClearingRegistry())) {
log.trace("stmt.id={} {}.id={} TCR is empty, skipping gateway request",
stmt.getId(), asts.a__t().getRegistryCode(), asts.a__t().getId());
} else {
} else if (TextUtil.isEmpty(sdf57.getSpecif())
|| !sdf57.getSpecif().contains(SpecifFlag.CS_BLKD.getKey())) {
log.trace("sending gateway request for stmt.id={}", stmt.getId());
Optional<AssetOperationListRequest> gtwReq = gatewayRequest(stmt,
company,
@ -399,17 +397,10 @@ public class Sdf57Executor extends AbstractExecutor<SDf57> {
searchResult.getCompany().getId());
asts = createRegistryIfNeeded.apply(new StmtCmpAcc(stmt, searchResult.getCompany(), searchResult.getAccount()));
if (OperationStatus.Executed.equalsByKey(stmt.getOperationStatus())) {
Optional<BigDecimal> dmiBalance = dmiService.setOk(searchResult.getTcr().getId(),
dmiService.setOk(searchResult.getTcr().getId(),
CurrencyCode.RUB.getKey(),
sdf57.getDbfId().toString());
asts.ifPresent(trio -> {
if (dmiBalance.isPresent() && sessionIsPresent) {
BigDecimal planBalance = safeBD(trio.a__t().getPlanBalance());
trio.a__t().setPlanBalance(planBalance.add(dmiBalance.get()));
registryImdg.update(trio.a__t());
}
assets.process(trio.a__b(), trio.a__t(), trio.a__f(), BigDecimal.ZERO);
});
asts.ifPresent(trio -> assets.process(trio.a__b(), trio.a__t(), trio.a__f(), BigDecimal.ZERO));
}
} else if (InOutDirection.in.equals(IEnumKey.getEnumByKey(InOutDirection.class, stmt.getInOutDirection()))) {
log.debug("stmt.id={} comment='{}' error: {}. Operating through DMAU registry",
@ -443,17 +434,10 @@ public class Sdf57Executor extends AbstractExecutor<SDf57> {
}
} else if (OperationStatus.Executed.equalsByKey(stmt.getOperationStatus())) {
if (asts.isPresent()) {
Optional<BigDecimal> dmiBalance = dmiService.setOk(asts.get().a__t().getTradingClearingRegistryId(),
dmiService.setOk(asts.get().a__t().getTradingClearingRegistryId(),
CurrencyCode.RUB.getKey(),
sdf57.getDbfId().toString());
asts.ifPresent(trio -> {
if (dmiBalance.isPresent() && sessionIsPresent) {
BigDecimal planBalance = safeBD(trio.a__t().getPlanBalance());
trio.a__t().setPlanBalance(planBalance.add(dmiBalance.get()));
registryImdg.update(trio.a__t());
}
assets.process(trio.a__b(), trio.a__t(), trio.a__f(), BigDecimal.ZERO);
});
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(),

View file

@ -125,13 +125,12 @@ public class PaymentInstructionOutboundService {
log.debug("all checks passed, accCred.id={}, accDeb.id={}, addressee.id={}, sender.id={}, amount: {}",
accCred.getId(), accDeb.getId(), addressee.getId(), sender.getId(), amount);
String purpose = tcr != null ? "Возврат денежных средств, свободных от обязательств с ТКР " + tcr.getCode() + " ." : "Вывод средств.";
String purpose = tcr != null ? "Вывод средств по ТКР " + tcr.getCode() + " ." : "Вывод средств.";
if (payload.getPaymentPurpose() != null) {
String reqPmtPrpse = payload.getPaymentPurpose();
if (!(reqPmtPrpse.endsWith(".") || reqPmtPrpse.endsWith("!") || reqPmtPrpse.endsWith("?") || reqPmtPrpse.endsWith(";"))) {
reqPmtPrpse += ".";
purpose += " " + payload.getPaymentPurpose();
if (!(purpose.endsWith(".") || purpose.endsWith("!") || purpose.endsWith("?") || purpose.endsWith(";"))) {
purpose += ".";
}
purpose = reqPmtPrpse + " " + purpose;
}
PaymentInstruction pmt = PaymentInstructionBuilderV2.builder(imdgProvider)
.registry(new Registry())

View file

@ -209,7 +209,6 @@ public class AssetTBFProcessing {
ast.setTradingClearingRegistry(tcr);
ast.setRegistryStatus(RegistryStatus.OK.getKey());
ast.setBalanceDimension(BalanceDimension.PICS.getKey());
ast.setOpenBalance(BigDecimal.ZERO);
RegistryManager.zeroState(ast);
Registry asf = copy(ast, RegistryUnit.F);
Registry asb = copy(ast, RegistryUnit.B);
@ -244,7 +243,6 @@ public class AssetTBFProcessing {
ast.setTradingClearingRegistry(tcr);
ast.setRegistryStatus(RegistryStatus.OK.getKey());
ast.setBalanceDimension(BalanceDimension.MONY.getKey());
ast.setOpenBalance(BigDecimal.ZERO);
RegistryManager.zeroState(ast);
Registry asf = ast.clone();
asf.setRegistryUnit(RegistryUnit.F.getKey());

View file

@ -4,18 +4,17 @@ import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Component;
import ru.clearing.classes.statics.data.registry.Registry;
import ru.clearing.classes.statics.data.sdf.SDf57;
import ru.clearing.classes.statics.data.statement.Statement;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.session.stage.util.RegistryUtil;
import ru.spcex.platform.enumeration.*;
import ru.spcex.platform.enumeration.RegistryDesignation;
import ru.spcex.platform.enumeration.RegistryStatus;
import ru.spcex.platform.enumeration.RegistryTradingParams;
import ru.spcex.platform.enumeration.RegistryUnit;
import ru.spcex.platform.imdg.api.Imdg;
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 ru.spcex.platform.utils.log.ExceptionUtils;
import ru.spcex.platform.utils.number.BigDecimalUtil;
import java.math.BigDecimal;
import java.time.Instant;
@ -25,14 +24,10 @@ import java.util.Optional;
public class DmiService {
private final Logger log = LoggerFactory.getLogger(getClass());
private final Imdg<Registry> rgsImdg;
private final Imdg<SDf57> sDf57Imdg;
private final Imdg<Statement> statementImdg;
private final RegistryManager rgsMng;
public DmiService(ImdgProvider imdgProvider, RegistryManager rgsMng) {
this.rgsImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
this.sDf57Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf57, SDf57.class);
this.statementImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Statement, Statement.class);
this.rgsMng = rgsMng;
}
@ -91,33 +86,6 @@ public class DmiService {
findDmiByComment(tcrId, securitySymbol, comment).ifPresentOrElse(d__i -> {
d__i.setContract(contract);
d__i.setUpdated(Instant.now());
try {
Long dbfId = Long.valueOf(contract);
SDf57 sdf57 = sDf57Imdg.getFirstObjectBySQL("dbfId=%d".formatted(dbfId));
if (sdf57 != null) {
ImdgPredicateBuilder pb = statementImdg.predicateBuilder();
ImdgPredicate stmtPrdct = pb.and(
pb.equals("inSDfId", sdf57.getId()),
pb.equals("inOutSDfType", InOutSDfType.type57.getKey())
);
Statement stmt = statementImdg.getFirstObjectByPredicate(stmtPrdct);
if (stmt != null && OperationStatus.Executed.equalsByKey(stmt.getOperationStatus())) {
d__i.setRegistryStatus(RegistryStatus.OK.getKey());
} else {
log.warn("tcrId={} securitySymbol={} comment={} sdf57 dbfId={} found, stmt.id={} status {}. d**i status will not be updated",
tcrId,
securitySymbol,
comment,
dbfId,
stmt != null ? stmt.getId() : null,
stmt != null ? stmt.getOperationStatus() : null);
}
}
} catch (Throwable e) {
log.error("sdf57 search by {}: error {}",
contract,
ExceptionUtils.getStackTrace(e));
}
rgsImdg.update(d__i);
log.trace("set {}.id={} contract {}", d__i.getRegistryCode(), d__i.getId(), contract);
}, () -> log.trace("didn't find DM*I by tcr.id={} securitySymbol={} comment={}",
@ -129,21 +97,18 @@ public class DmiService {
* ищем D**I по tcrId/security/contract
* задаем статус ОК
*/
public Optional<BigDecimal> setOk(Long tcrId,
public void setOk(Long tcrId,
String securitySymbol,
String contract) {
BigDecimal[] dmiBalance = {null};
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());
rgsImdg.update(d__i);
dmiBalance[0] = BigDecimalUtil.safeBD(d__i.getBalance());
}, () -> {
log.trace("didn't find DM*I by tcr.id={} securitySymbol={} contract={}",
tcrId, securitySymbol, contract);
});
return Optional.ofNullable(dmiBalance[0]);
}
private Optional<Registry> findDmiByContract(Long tcrId, String securitySymbol, String contract) {

View file

@ -38,17 +38,11 @@ public class RegistryManager {
public Optional<Registry> searchDmxByCounterParty(Registry rgs) {
if (TextUtil.isEmpty(rgs.getContract())) return Optional.empty();
ImdgPredicateBuilder pb = rgsImdg.predicateBuilder();
ImdgPredicate prdct = pb.and(
pb.sql(RegistryCodeSqlBuilder.getInstance(DM_X).build()),
pb.equals("companyId", rgs.getCounterPartyId()),
pb.equals("contract", rgs.getContract()),
pb.or(
pb.equals("registryStatus", RegistryStatus.PROC.getKey()),
pb.equals("registryStatus", RegistryStatus.MNG.getKey())
)
);
Registry dmx = rgsImdg.getFirstObjectByPredicate(prdct);
String sqlCondition = String.format("(%s) and companyId = %d and contract = '%s'",
RegistryCodeSqlBuilder.getInstance(DM_X).build(),
rgs.getCounterPartyId(),
rgs.getContract());
Registry dmx = rgsImdg.getFirstObjectBySQL(sqlCondition);
return Optional.ofNullable(dmx);
}

View file

@ -131,7 +131,7 @@ public enum Sdf06NewValidationRule implements IValidationRule<ImdgValidationCont
prdctBldr.and(
prdctByRegistryCode.apply(RegistryTradingParams.DM_T),
prdctBldr.or(
prdctBldr.equals("registryStatus", RegistryStatus.PROC.getKey()),
prdctBldr.equals("registryStatus", RegistryStatus.OK.getKey()),
prdctBldr.equals("registryStatus", RegistryStatus.MNG.getKey())
)
)

View file

@ -63,7 +63,7 @@ public enum StatusExtractValidationRule implements IValidationRule<ImdgValidatio
return empty();
}
},
RequestStatusValid() {
RegistryStatusValid() {
@Override
public Optional<EnumMessage> validate(ImdgValidationContext<RegistryChangeStatusExtractRequest> context) {
RegistryChangeStatusExtractRequest validatedObject = context.getValidatedObject();
@ -76,17 +76,6 @@ public enum StatusExtractValidationRule implements IValidationRule<ImdgValidatio
}
return empty();
}
},
RegistryStatusValid() {
@Override
public Optional<EnumMessage> validate(ImdgValidationContext<RegistryChangeStatusExtractRequest> context) {
Registry rgs = context.getStoredObject(Stored.PresentById);
RegistryStatus status = IEnumKey.getEnumByKey(RegistryStatus.class, rgs.getRegistryStatus());
if (RegistryStatus.OK.equals(status)) {
return of(ClearingError.ObligationsAlreadyCalculated);
}
return empty();
}
}
;
private final static Logger log = LoggerFactory.getLogger(StatusExtractValidationRule.class);

View file

@ -0,0 +1,57 @@
package ru.spcex.clearing.messaging;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.platform.messaging.service.IKafkaRequestStatusProcessor;
import ru.spcex.clearing.platform.messaging.service.RequestInfo;
import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate;
import ru.spcex.clearing.platform.messaging.service.Status;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.utils.log.ExceptionUtils;
public class KafkaRequestStatusProcessorImpl implements IKafkaRequestStatusProcessor {
private final Logger log = LoggerFactory.getLogger(getClass());
private final Imdg<RequestInfo> rqstInfImdg;
public KafkaRequestStatusProcessorImpl(ImdgProvider imdgProvider) {
this.rqstInfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class);
}
@Override
public void updateRequestStatus(BaseRequest<?> o, Object response) {
try {
RequestInfoUpdate updateEvent;
if (response == null) {
//default response
updateEvent = new RequestInfoUpdate();
updateEvent.setId(o.getId());
updateEvent.setStatus(Status.Success);
} else {
updateEvent = (RequestInfoUpdate) response;
}
RequestInfo rqstInfo = rqstInfImdg.getSingleObjectByID(o.getId());
if (rqstInfo == null) {
log.trace("cannot find requestInfo.id={} for update '{}'", o.getId(), updateEvent.getStatus());
return;
}
rqstInfo.setStatus(updateEvent.getStatus());
rqstInfo.setMessage(updateEvent.getMessage());
rqstInfImdg.update(rqstInfo);
} catch (Exception e) {
log.error(ExceptionUtils.getStackTrace(e));
}
}
@Override
public void updateRequestStatusOnError(BaseRequest<?> o, Object response) {
RequestInfoUpdate updateEvent;
//default response
updateEvent = new RequestInfoUpdate();
updateEvent.setId(o.getId());
updateEvent.setStatus(Status.Error);
updateRequestStatus(o, updateEvent);
}
}

View file

@ -6,7 +6,7 @@
/* Business objects */
INSERT INTO COMPANY(ID, FULL_NAME, SHORT_NAME, TRADING_CODE, CLEARING_CODE, WORKFLOW_STATUS) values (1, 'АО Санкт-Петербургская Валютная Биржа', 'СПВБ', '001', '001', 'ACTV') ON CONFLICT (ID) DO UPDATE SET FULL_NAME = EXCLUDED.FULL_NAME, SHORT_NAME = EXCLUDED.SHORT_NAME, TRADING_CODE = EXCLUDED.TRADING_CODE, CLEARING_CODE = EXCLUDED.CLEARING_CODE, WORKFLOW_STATUS = EXCLUDED.WORKFLOW_STATUS;
INSERT INTO COMPANY(ID, FULL_NAME, SHORT_NAME, TRADING_CODE, CLEARING_CODE, WORKFLOW_STATUS) values (1, 'АО Санкт-Петербургская Валютная Биржа', 'СПВБ', '', '', 'ACTV') ON CONFLICT (ID) DO UPDATE SET FULL_NAME = EXCLUDED.FULL_NAME, SHORT_NAME = EXCLUDED.SHORT_NAME, TRADING_CODE = EXCLUDED.TRADING_CODE, CLEARING_CODE = EXCLUDED.CLEARING_CODE, WORKFLOW_STATUS = EXCLUDED.WORKFLOW_STATUS;
INSERT INTO COMPANY(ID, FULL_NAME, SHORT_NAME, TRADING_CODE, CLEARING_CODE, WORKFLOW_STATUS) values (2, 'ЗАО «Петербургский Расчетный Центр»', 'ПРЦ', '', '', 'ACTV') ON CONFLICT (ID) DO UPDATE SET FULL_NAME = EXCLUDED.FULL_NAME, SHORT_NAME = EXCLUDED.SHORT_NAME, TRADING_CODE = EXCLUDED.TRADING_CODE, CLEARING_CODE = EXCLUDED.CLEARING_CODE, WORKFLOW_STATUS = EXCLUDED.WORKFLOW_STATUS;
@ -26,7 +26,7 @@ INSERT INTO ACCOUNT(ID, COMPANY_ID, ACCOUNT, ACCOUNT_TYPE, STATUS, PROCESSING_SI
INSERT INTO ACCOUNT(ID, COMPANY_ID, ACCOUNT, ACCOUNT_TYPE, STATUS, PROCESSING_SIGN) values (2, 1, '30414810600000007000', 'ANLT', 'ACTV', 'ALWD') ON CONFLICT (ID) DO UPDATE SET COMPANY_ID = EXCLUDED.COMPANY_ID, ACCOUNT = EXCLUDED.ACCOUNT, ACCOUNT_TYPE = EXCLUDED.ACCOUNT_TYPE, STATUS = EXCLUDED.STATUS, PROCESSING_SIGN = EXCLUDED.PROCESSING_SIGN;
INSERT INTO ACCOUNT(ID, COMPANY_ID, ACCOUNT, ACCOUNT_TYPE, STATUS, PROCESSING_SIGN) values (3, 1, 'SPCEX', 'DTRN', 'ACTV', 'ALWD') ON CONFLICT (ID) DO UPDATE SET COMPANY_ID = EXCLUDED.COMPANY_ID, ACCOUNT = EXCLUDED.ACCOUNT, ACCOUNT_TYPE = EXCLUDED.ACCOUNT_TYPE, STATUS = EXCLUDED.STATUS, PROCESSING_SIGN = EXCLUDED.PROCESSING_SIGN;
INSERT INTO ACCOUNT(ID, COMPANY_ID, ACCOUNT, ACCOUNT_TYPE, STATUS, PROCESSING_SIGN) values (3, 1, '700100000AT0', 'DTRN', 'ACTV', 'ALWD') ON CONFLICT (ID) DO UPDATE SET COMPANY_ID = EXCLUDED.COMPANY_ID, ACCOUNT = EXCLUDED.ACCOUNT, ACCOUNT_TYPE = EXCLUDED.ACCOUNT_TYPE, STATUS = EXCLUDED.STATUS, PROCESSING_SIGN = EXCLUDED.PROCESSING_SIGN;
INSERT INTO MARKET(ID, MARKET_TYPE, SECTION, SETTLEMENT_CURRENCY, CODE, NAME, EXCHANGE_ID, DESCRIPTION) values (1, 'SCND', 'FOND', 'RUB', 'UESC', 'Режим непрерывных торгов', 1, 'Режим непрерывных торгов: Обыкновенные акции') ON CONFLICT (ID) DO UPDATE SET MARKET_TYPE = EXCLUDED.MARKET_TYPE, SECTION = EXCLUDED.SECTION, SETTLEMENT_CURRENCY = EXCLUDED.SETTLEMENT_CURRENCY, CODE = EXCLUDED.CODE, NAME = EXCLUDED.NAME, EXCHANGE_ID = EXCLUDED.EXCHANGE_ID, DESCRIPTION = EXCLUDED.DESCRIPTION;

View file

@ -18,22 +18,22 @@ public abstract class TemplateEventMapStore<T extends BusinessEvent<? extends Sp
}
@Override
public Iterable<Long> loadAllKeys() {
public final Iterable<Long> loadAllKeys() {
return Collections.emptyList();
}
@Override
public T load(Long id) {
public final T load(Long id) {
return null;
}
@Override
public Collection<T> load(Collection<Long> keys) {
public final Collection<T> load(Collection<Long> keys) {
return Collections.emptyList();
}
@Override
public Map<Long, T> loadAll(Collection<Long> keys) {
public final Map<Long, T> loadAll(Collection<Long> keys) {
return Collections.emptyMap();
}

View file

@ -1,6 +1,5 @@
package ru.spcex.clearing.imdg.businessevent;
import org.springframework.dao.support.DataAccessUtils;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Component;
import ru.clearing.classes.statics.data.registry.Registry;
@ -9,11 +8,6 @@ import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.imdg.base.TemplateEventMapStore;
import ru.spcex.platform.utils.time.TimeUtil;
import java.math.BigDecimal;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.util.*;
@Component
public class RegistryHistoryMapStore extends TemplateEventMapStore<RegistryHistory> {
@ -21,154 +15,23 @@ public class RegistryHistoryMapStore extends TemplateEventMapStore<RegistryHisto
super(jdbcTemplate);
}
@Override
public String getTableName() {
return "REGISTRY_HISTORY";
}
@Override
public String getMapName() {
return IMDGDistributedNames.Map_RegistryHistory;
}
@Override
public String getTableName() {
return "REGISTRY_HISTORY";
}
@Override
public String[] getFields() {
return new String[]{
"ID", "EVENT_TIME", "EVENT_USER_ID", "EVENT_TYPE", "REGISTRY_ID", "CREATED_AT", "UPDATED_AT",
"COMPANY_ID", "TRADING_CODE", "CLEARING_CODE", "SHORT_NAME", "FULL_NAME", "ACCOUNT_ID",
"ACCOUNT_TYPE", "ACCOUNT", "REGISTRY_DESIGNATION", "REGISTRY_INSTRUMENT_TYPE", "REGISTRY_CAPACITY",
"REGISTRY_UNIT", "REGISTRY_CODE", "TRADING_CLEARING_REGISTRY_ID", "TRADING_CLEARING_REGISTRY",
"REGISTRY_STATUS", "SECURITY_ID", "SECURITY_SYMBOL", "BALANCE", "OPEN_BALANCE", "CLOSE_BALANCE",
"CREDIT", "DEBIT", "SETTLED_CREDIT", "SETTLED_DEBIT", "CHECK_BALANCE", "DIFF_BALANCE", "PLAN_BALANCE",
"BALANCE_DIMENSION", "SETTLEMENT_DATE", "SETTLEMENT_CODE", "TRADING_DATE", "CLEARING_DATE", "REFUND_DATE",
"VALUE_DATE", "PRICE", "CONTRACT", "COUNTER_PARTY_ID", "COMMENT", "PARENT_ID", "GROUP_ID", "SESSION_ID",
"SESSION_TYPE", "PAYMENT_ID"
return new String[]{"ID", "EVENT_TIME", "EVENT_USER_ID", "EVENT_TYPE",
"REGISTRY_ID", "CREATED_AT", "UPDATED_AT", "COMPANY_ID", "TRADING_CODE", "CLEARING_CODE", "SHORT_NAME", "FULL_NAME", "ACCOUNT_ID", "ACCOUNT_TYPE", "ACCOUNT", "REGISTRY_DESIGNATION", "REGISTRY_INSTRUMENT_TYPE", "REGISTRY_CAPACITY", "REGISTRY_UNIT", "REGISTRY_CODE", "TRADING_CLEARING_REGISTRY_ID", "TRADING_CLEARING_REGISTRY", "REGISTRY_STATUS", "SECURITY_ID", "SECURITY_SYMBOL", "BALANCE", "OPEN_BALANCE", "CLOSE_BALANCE", "CREDIT", "DEBIT", "SETTLED_CREDIT", "SETTLED_DEBIT", "CHECK_BALANCE", "DIFF_BALANCE", "PLAN_BALANCE", "BALANCE_DIMENSION", "SETTLEMENT_DATE", "SETTLEMENT_CODE", "TRADING_DATE", "CLEARING_DATE", "REFUND_DATE", "VALUE_DATE", "PRICE", "CONTRACT", "COUNTER_PARTY_ID", "COMMENT", "PARENT_ID", "GROUP_ID", "SESSION_ID", "SESSION_TYPE", "PAYMENT_ID"
};
}
@Override
public Iterable<Long> loadAllKeys() {
return defaultLoadAllKeysOnTodayByField("EVENT_TIME", true);
}
@Override
public RegistryHistory load(Long id) {
if (!isLoadable(id)) return null;
Collection<RegistryHistory> rows;
List<Long> list = Collections.singletonList(id);
try {
rows = load(list);
} catch (Throwable e) { // one retry
rows = load(list);
}
RegistryHistory obj = DataAccessUtils.singleResult(rows);
if (obj != null && isLoadable(obj))
return obj;
else
return null;
}
@Override
public Map<Long, RegistryHistory> loadAll(Collection<Long> keys) {
log.debug("loadAll from " + getTableName() + " " + keys.size() + " keys");
Map<Long, RegistryHistory> result = new HashMap<>();
long start = System.currentTimeMillis();
// загрузить данные по ключам частями, чтобы не выйти за ограничения базы по кол-ву элементов в in clause
List<Long> keysSubList = new ArrayList<>(MAX_IN_CLAUSE_SIZE);
for (Iterator<Long> iterator = keys.iterator(); iterator.hasNext(); ) {
Long key = iterator.next();
keysSubList.add(key);
if (keysSubList.size() == MAX_IN_CLAUSE_SIZE || !iterator.hasNext()) {
Collection<RegistryHistory> rows = load(keysSubList);
for (RegistryHistory row : rows) {
result.put(row.getId(), row);
}
keysSubList.clear();
}
}
log.debug("loadAll from " + getTableName() + " " + keys.size() + " keys done in " + (System.currentTimeMillis() - start) + "ms");
return result;
}
@Override
public Collection<RegistryHistory> load(Collection<Long> keys) {
Map<String, Collection<Long>> paramMap = Collections.singletonMap("ids", keys);
return namedParameterJdbcTemplate.query("select * from " + getTableName() + " where id in (:ids)", paramMap,
(resultSet, i) -> objectReader(resultSet));
}
protected RegistryHistory objectReader(ResultSet resultSet) throws SQLException {
RegistryHistory registryHistory = new RegistryHistory();
Registry object = new Registry();
registryHistory.setId(resultSet.getObject("ID", Long.class));
registryHistory.setEventTime(getInstantFromTimestamp(resultSet, "EVENT_TIME"));
registryHistory.setUserId(resultSet.getObject("EVENT_USER_ID", Long.class));
registryHistory.setEventType(resultSet.getObject("EVENT_TYPE", String.class));
object.setCreated(getInstantFromTimestamp(resultSet, "CREATED_AT"));
object.setUpdated(getInstantFromTimestamp(resultSet, "UPDATED_AT"));
object.setId(resultSet.getObject("REGISTRY_ID", Long.class));
object.setCompanyId(resultSet.getObject("COMPANY_ID", Long.class));
object.setTradingCode(resultSet.getObject("TRADING_CODE", String.class));
object.setClearingCode(resultSet.getObject("CLEARING_CODE", String.class));
object.setShortName(resultSet.getObject("SHORT_NAME", String.class));
object.setFullName(resultSet.getObject("FULL_NAME", String.class));
object.setAccountId(resultSet.getObject("ACCOUNT_ID", Long.class));
object.setAccountType(resultSet.getObject("ACCOUNT_TYPE", String.class));
object.setAccount(resultSet.getObject("ACCOUNT", String.class));
object.setRegistryDesignation(resultSet.getObject("REGISTRY_DESIGNATION", String.class));
object.setRegistryInstrumentType(resultSet.getObject("REGISTRY_INSTRUMENT_TYPE", String.class));
object.setRegistryCapacity(resultSet.getObject("REGISTRY_CAPACITY", String.class));
object.setRegistryUnit(resultSet.getObject("REGISTRY_UNIT", String.class));
object.setRegistryCode(resultSet.getObject("REGISTRY_CODE", String.class));
object.setTradingClearingRegistryId(resultSet.getObject("TRADING_CLEARING_REGISTRY_ID", Long.class));
object.setTradingClearingRegistry(resultSet.getObject("TRADING_CLEARING_REGISTRY", String.class));
object.setRegistryStatus(resultSet.getObject("REGISTRY_STATUS", String.class));
object.setSecurityId(resultSet.getObject("SECURITY_ID", Long.class));
object.setSecuritySymbol(resultSet.getObject("SECURITY_SYMBOL", String.class));
object.setBalance(resultSet.getObject("BALANCE", BigDecimal.class));
object.setOpenBalance(resultSet.getObject("OPEN_BALANCE", BigDecimal.class));
object.setCloseBalance(resultSet.getObject("CLOSE_BALANCE", BigDecimal.class));
object.setCredit(resultSet.getObject("CREDIT", BigDecimal.class));
object.setDebit(resultSet.getObject("DEBIT", BigDecimal.class));
object.setSettledCredit(resultSet.getObject("SETTLED_CREDIT", BigDecimal.class));
object.setSettledDebit(resultSet.getObject("SETTLED_DEBIT", BigDecimal.class));
object.setCheckBalance(resultSet.getObject("CHECK_BALANCE", BigDecimal.class));
object.setDiffBalance(resultSet.getObject("DIFF_BALANCE", BigDecimal.class));
object.setPlanBalance(resultSet.getObject("PLAN_BALANCE", BigDecimal.class));
object.setBalanceDimension(resultSet.getObject("BALANCE_DIMENSION", String.class));
object.setSettlementDate(getLocalDateFromSqlDate(resultSet, "SETTLEMENT_DATE"));
object.setSettlementCode(resultSet.getObject("SETTLEMENT_CODE", String.class));
object.setTradingDate(getLocalDateFromSqlDate(resultSet, "TRADING_DATE"));
object.setClearingDate(getLocalDateFromSqlDate(resultSet, "CLEARING_DATE"));
object.setRefundDate(getLocalDateFromSqlDate(resultSet, "REFUND_DATE"));
object.setValueDate(getLocalDateFromSqlDate(resultSet, "VALUE_DATE"));
object.setPrice(resultSet.getObject("PRICE", BigDecimal.class));
object.setContract(resultSet.getObject("CONTRACT", String.class));
object.setCounterPartyId(resultSet.getObject("COUNTER_PARTY_ID", Long.class));
object.setComment(resultSet.getObject("COMMENT", String.class));
object.setParentId(resultSet.getObject("PARENT_ID", Long.class));
object.setGroupId(resultSet.getObject("GROUP_ID", Long.class));
object.setSessionId(resultSet.getObject("SESSION_ID", Long.class));
object.setSessionType(resultSet.getObject("SESSION_TYPE", String.class));
object.setPaymentId(resultSet.getObject("PAYMENT_ID", Long.class));
registryHistory.setObject(object);
return registryHistory;
}
@Override
public Object[] objectToField(RegistryHistory historyLog) {
Registry object = historyLog.getObject();

View file

@ -41,9 +41,9 @@ public class DbConnectionConfig {
cpds.setJdbcUrl(dbPath);
cpds.setUser(login);
cpds.setPassword(password);
cpds.setInitialPoolSize(settings.getMinPoolSize());
cpds.setMinPoolSize(settings.getMinPoolSize());
cpds.setMaxPoolSize(settings.getMaxPoolSize());
cpds.setInitialPoolSize(10);
cpds.setMinPoolSize(10);
cpds.setMaxPoolSize(30);
int numHelperThreads = Runtime.getRuntime().availableProcessors() * 2;
cpds.setNumHelperThreads(numHelperThreads);
cpds.setCheckoutTimeout(timeoutSec * 1000);

View file

@ -32,12 +32,11 @@ public class HazelcastConfiguration {
@Bean
public Config hazelCastConfig() {
Config config = new Config();
config.setInstanceName(hzSettings.getInstanceName());
config.setInstanceName("instance");
config.setGroupConfig(new GroupConfig()
.setName(hzSettings.getLogin())
.setPassword(hzSettings.getPassword())
);
config.setManagementCenterConfig(new ManagementCenterConfig().setEnabled(false));
config.setProperty("hazelcast.shutdownhook.enabled", "true");
config.setProperty("hazelcast.logging.type", "slf4j");
config.setProperty("hazelcast.operation.call.timeout.millis", "600000");

View file

@ -49,7 +49,7 @@ public class PoolMapConfigs implements InitializingBean {
private MapStoreConfig makeDefaultMapStoreConfig(MapLoader<Long, ?> mapBean) {
return new MapStoreConfig()
.setImplementation(mapBean)
.setWriteDelaySeconds(settings.getHazelcast().getDbSyncSeconds());
;//todo config: .setWriteDelaySeconds(configRoot.getIMDG().getDbSyncSeconds());
}
public ScheduledExecutorConfig makeDefaultScheduledExecutorConfig(String name) {

View file

@ -4,8 +4,6 @@ public class DatabaseSettings {
private String login;
private String password;
private String url;
private int minPoolSize = 10;
private int maxPoolSize = 30;
public String getLogin() {
return login;
@ -30,20 +28,4 @@ public class DatabaseSettings {
public void setUrl(String url) {
this.url = url;
}
public int getMaxPoolSize() {
return maxPoolSize;
}
public void setMaxPoolSize(int maxPoolSize) {
this.maxPoolSize = maxPoolSize;
}
public int getMinPoolSize() {
return minPoolSize;
}
public void setMinPoolSize(int minPoolSize) {
this.minPoolSize = minPoolSize;
}
}

View file

@ -9,8 +9,6 @@ public class HazelcastServerSettings {
private String password;
private List<String> clusterMembers;
private Integer backupCount;
private String instanceName = "instance";
private Integer dbSyncSeconds = 10;
public int getListenPort() {
return listenPort;
@ -51,20 +49,4 @@ public class HazelcastServerSettings {
public void setBackupCount(Integer backupCount) {
this.backupCount = backupCount;
}
public String getInstanceName() {
return instanceName;
}
public void setInstanceName(String instanceName) {
this.instanceName = instanceName;
}
public Integer getDbSyncSeconds() {
return dbSyncSeconds;
}
public void setDbSyncSeconds(Integer dbSyncSeconds) {
this.dbSyncSeconds = dbSyncSeconds;
}
}

View file

@ -3,11 +3,7 @@ imdg.hazelcast.login=dev
imdg.hazelcast.password=dev-pass
imdg.hazelcast.cluster-members[0]=127.0.0.1
imdg.hazelcast.backup-count=3
imdg.hazelcast.instance-name=instance
imdg.hazelcast.db-sync-seconds=10
imdg.database.login=clearing
imdg.database.password=Aa111111
imdg.database.url=jdbc:postgresql://10.200.200.133:5432/clearing?currentSchema=clearing_prod
#imdg.database.url=jdbc:postgresql://10.200.200.133:5432/postgres?currentSchema=clearing_tester
imdg.database.min-pool-size=10
imdg.database.max-pool-size=30
#imdg.database.url=jdbc:postgresql://10.200.200.133:5432/postgres?currentSchema=clearing_tester

View file

@ -15,13 +15,14 @@ import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
import ru.spcex.platform.imdg.api.predicate.specific.RegistryCodeSqlBuilder;
import java.math.BigDecimal;
import java.time.LocalDate;
import java.util.ArrayList;
import java.util.Collection;
import java.util.List;
import java.util.function.Function;
import java.util.stream.Collectors;
import static ru.spcex.platform.enumeration.RegistryTradingParams.AM_F;
import static ru.spcex.platform.enumeration.RegistryTradingParams.DM_T;
import static ru.spcex.platform.enumeration.RegistryTradingParams.*;
@Service
public class MoneyExporterService extends AbstractExporterService {
@ -51,18 +52,29 @@ public class MoneyExporterService extends AbstractExporterService {
log.debug("Selected {} registriesA.", registriesA.size());
ImdgPredicateBuilder pb = registryImdg.predicateBuilder();
LocalDate today = LocalDate.now();
Function<Registry, ImdgPredicate> prdctDM_T = registry -> pb.and(
pb.sql(RegistryCodeSqlBuilder.getInstance(DM_T).build()),
pb.equals("registryStatus", RegistryStatus.PROC.getKey()),
pb.equals("settlementDate", today),
pb.equals("tradingClearingRegistry", registry.getTradingClearingRegistry()),
pb.equals("securityId", registry.getSecurityId())
);
Function<Registry, ImdgPredicate> prdctDM_I = registry -> pb.and(
pb.sql(RegistryCodeSqlBuilder.getInstance(DM_I).build()),
pb.equals("tradingClearingRegistry", registry.getTradingClearingRegistry()),
pb.equals("securityId", registry.getSecurityId()),
pb.not(pb.equals("registryStatus", RegistryStatus.OK.getKey()))
);
for (Registry registry : registriesA) {
if (checkNotBlocked(registry)) {
BigDecimal sumPositive = BigDecimal.ZERO;
if (registry.getTradingClearingRegistry() != null && registry.getSecurityId() != null) {
Collection<Registry> regsDMT = registryImdg.getCollectionObjectsByPredicate(prdctDM_T.apply(registry));
Collection<Registry> regsDMI = registryImdg.getCollectionObjectsByPredicate(prdctDM_I.apply(registry));
log.trace("Found {} DM_I registries for registry AM_F {}", regsDMI.size(), registry);
Collection<Registry> regsDMT = filterByDate(registryImdg.getCollectionObjectsByPredicate(prdctDM_T.apply(registry)));
log.trace("Found {} DM_T registries for registry AM_F {}", regsDMT.size(), registry);
sumPositive = sumNegativeABS(regsDMI);
sumPositive = sumPositive.add(sum(regsDMT));
}
limFileRows.add(getRow(registry, sumPositive));
@ -72,6 +84,12 @@ public class MoneyExporterService extends AbstractExporterService {
return limFileRows;
}
List<Registry> filterByDate(Collection<Registry> regs) {
return regs.stream()
.filter((Registry reg) -> reg.getSettlementDate().equals(reg.getClearingDate()))
.collect(Collectors.toList());
}
public String getRow(Registry registryA, BigDecimal sumPositive) {
StringBuilder row = new StringBuilder();

View file

@ -7,8 +7,10 @@ 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;
@ -16,8 +18,10 @@ 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.AS_F;
import static ru.spcex.platform.enumeration.RegistryTradingParams.DS_T;
@Service
public class SecurityExporterService extends AbstractExporterService {
@ -45,6 +49,13 @@ public class SecurityExporterService extends AbstractExporterService {
log.debug("Started loading and formation of DEPO file lines");
ImdgPredicateBuilder pb = registryImdg.predicateBuilder();
Function<Registry, ImdgPredicate> prdct = registry -> pb.and(
pb.sql(RegistryCodeSqlBuilder.getInstance(DS_T).build()),
pb.equals("registryStatus", RegistryStatus.PROC.getKey()),
pb.equals("companyId", registry.getCompanyId()),
pb.equals("tradingClearingRegistry", registry.getTradingClearingRegistry()),
pb.equals("securityId", registry.getSecurityId())
);
List<String> limFileRows = new ArrayList<>();
Collection<Registry> registriesA = registryImdg.getCollectionObjectsBySQL(RegistryCodeSqlBuilder.getInstance(AS_F).build());
@ -52,14 +63,18 @@ public class SecurityExporterService extends AbstractExporterService {
for (Registry registry : registriesA) {
if (checkNotBlocked(registry)) {
limFileRows.add(getRow(registry));
BigDecimal sumNegative = BigDecimal.ZERO;
if (registry.getCompanyId() != null && registry.getTradingClearingRegistry() != null && registry.getSecurityId() != null) {
sumNegative = sumNegativeABS(registryImdg.getCollectionObjectsByPredicate(prdct.apply(registry)));
}
limFileRows.add(getRow(registry, sumNegative));
}
}
log.debug("Successfully completed the formation of rows: {} for export DEPO", limFileRows.size());
return limFileRows;
}
public String getRow(Registry registry) {
public String getRow(Registry registry, BigDecimal sumPositive) {
StringBuilder row = new StringBuilder();
row.append("DEPO: FIRM_ID = ");
@ -72,7 +87,7 @@ public class SecurityExporterService extends AbstractExporterService {
row.append(registry.getTradingClearingRegistry());
row.append("; OPEN_BALANCE = ");
BigDecimal balance = registry.getBalance() != null ? registry.getBalance() : BigDecimal.ZERO;
BigDecimal balance = registry.getBalance() != null ? registry.getBalance().subtract(sumPositive) : BigDecimal.ZERO;
row.append(LimFormat.toStringD2(balance));
row.append("; OPEN_LIMIT = 0");

View file

@ -9,9 +9,11 @@ import ru.spcex.platform.enumeration.RegistryUnit;
import ru.spcex.platform.enumeration.ServiceStatus;
import javax.annotation.PostConstruct;
import java.math.BigDecimal;
import java.util.Collection;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static ru.spcex.clearing.test.TestUtils.clearAllInImdg;
class SecurityExporterServiceTest extends AbstractServiceTest {
@ -48,7 +50,7 @@ class SecurityExporterServiceTest extends AbstractServiceTest {
tradingClearingRegistryImdg.insert(tradingClearingRegistry);
limFileRows = securityExporterService.getLimFileRows();
// assertEquals(2, limFileRows.size());
// assertTrue(limFileRows.contains(securityExporterService.getRow(registryA, new BigDecimal("5.00"))));
assertEquals(2, limFileRows.size());
assertTrue(limFileRows.contains(securityExporterService.getRow(registryA, new BigDecimal("5.00"))));
}
}

View file

@ -9,7 +9,6 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.register.DepoBalanceRegister;
import ru.clearing.classes.statics.data.registry.Registry;
import ru.clearing.classes.statics.data.registry.RegistryHistory;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
@ -22,30 +21,34 @@ 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 ru.spcex.platform.utils.enumeration.IEnumKey;
import ru.spcex.platform.utils.validation.IValidator;
import java.math.BigDecimal;
import java.time.Instant;
import java.time.temporal.ChronoUnit;
import java.util.Collection;
import java.util.Map;
import java.util.Optional;
import java.util.function.Function;
import static ru.spcex.platform.enumeration.Task.createRegistry_GRRT;
@Service
public class DepoBalanceRegisterService extends QueueConsumer implements InitializingBean {
public class DepoBalanceRegisterService extends QueueConsumer implements InitializingBean, ICheckDuplicate<Registry, DepoBalanceRegister> {
private final Logger log = LoggerFactory.getLogger(getClass());
private final Imdg<DepoBalanceRegister> depoBalanceRegisterMap;
private final Imdg<RegistryHistory> registryHistoryMap;
private final Imdg<Registry> registryMap;
private final Function<Map<String, ?>, IValidator> fieldValuesValidator;
private final PreClearMap<DepoBalanceRegister> preClearMap;
@Autowired
public DepoBalanceRegisterService(Consumer<String, Object> kafkaQueue,
Producer<String, Object> kafkaProducer,
ImdgProvider imdgProvider) {
ImdgProvider imdgProvider,
Function<Map<String, ?>, IValidator> fieldValuesValidator) {
super(kafkaQueue, kafkaProducer);
this.depoBalanceRegisterMap = imdgProvider.getImdg(IMDGDistributedNames.Map_DepoBalanceRegister, DepoBalanceRegister.class);
this.registryHistoryMap = imdgProvider.getImdg(IMDGDistributedNames.Map_RegistryHistory, RegistryHistory.class);
this.registryMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
this.fieldValuesValidator = fieldValuesValidator;
preClearMap = PreClearMap.instanceForInstantField(depoBalanceRegisterMap, "created");
}
@ -59,32 +62,35 @@ public class DepoBalanceRegisterService extends QueueConsumer implements Initial
public void depoBalanceRegisterNew(BaseRequest<LauncherCommandRequest> userRequest) {
log.debug("LauncherCommandRequest received from {}", createRegistry_GRRT.topic());
ImdgPredicate predicate = getPredicateForRegistries(userRequest.getRequestPayload().getCompanyId());
Collection<RegistryHistory> registryHistories = registryHistoryMap.getCollectionObjectsByPredicate(predicate);
log.trace("Started searching depoBalanceRegister in register by '{}'", predicate);
if (registryHistories.isEmpty()) {
log.debug("No registry with such conditions '{}'", predicate);
ImdgPredicate prdctForRegistries = getPredicateForRegistries(userRequest.getRequestPayload().getCompanyId());
Collection<Registry> registries = registryMap.getCollectionObjectsByPredicate(prdctForRegistries);
log.trace("Started searching depoBalanceRegister in register by '{}'", prdctForRegistries);
if (registries.isEmpty()) {
log.debug("No registry with such conditions '{}'", prdctForRegistries);
return;
}
log.debug("Select {} Registry history by query {}", registryHistories.size(), predicate);
log.debug("Select {} Registry by query {}", registries.size(), prdctForRegistries);
preClearMap.preClearMap();
for (RegistryHistory registryHistory : registryHistories) {
if (BigDecimal.ZERO.compareTo(registryHistory.getObject().getDiffBalance()) != 0) continue;
insertDepoBalanceRegister(registryHistory);
registries.forEach((registry -> {
if (isDuplicateInMap(registry).isEmpty()) {
insertDepoBalanceRegister(registry);
}
}
));
log.debug("Successfully processed");
}
private void insertDepoBalanceRegister(RegistryHistory registryHistory) {
private void insertDepoBalanceRegister(Registry registry) {
log.trace("Started generating DepoBalanceRegister entity...");
Registry registry = registryHistory.getObject();
DepoBalanceRegister depoBalanceRegister = new DepoBalanceRegister();
depoBalanceRegister.setUpdated(Instant.now());
depoBalanceRegister.setCreated(depoBalanceRegister.getUpdated());
depoBalanceRegister.setCompanyId(registry.getCompanyId());
AccountType accountType = IEnumKey.getEnumByKey(AccountType.class, registry.getAccountType());
if (AccountType.Depo == accountType) {
depoBalanceRegister.setSessionId(registry.getSessionId());
if (AccountType.Depo.equalsByKey(registry.getAccountType())) {
depoBalanceRegister.setDepoCode(registry.getAccount());
} else if (registry.getAccountType() != null) {
log.warn("Unexpected registry[{}].accountType={}", registry.getId(), registry.getAccountType());
}
depoBalanceRegister.setQuantity(registry.getBalance());
depoBalanceRegister.setSecuritySymbol(registry.getSecuritySymbol());
@ -93,29 +99,26 @@ public class DepoBalanceRegisterService extends QueueConsumer implements Initial
}
private ImdgPredicate getPredicateForRegistries(Long companyId) {
Instant now = Instant.now().truncatedTo(ChronoUnit.DAYS);
ImdgPredicateBuilder pb = registryHistoryMap.predicateBuilder();
RegistryCodeSqlBuilder codeSql = RegistryCodeSqlBuilder.getInstance(RegistryTradingParams.AS_T).prefix("object");
ImdgPredicate registryPredicate;
if (companyId == null) {
registryPredicate = pb.sql(codeSql.build());
} else {
registryPredicate = pb.and(
pb.equals("object.companyId", companyId),
ImdgPredicateBuilder pb = registryMap.predicateBuilder();
RegistryCodeSqlBuilder codeSql = RegistryCodeSqlBuilder.getInstance(RegistryTradingParams.AS_T);
if (companyId == null) // default use case
return pb.sql(codeSql.build());
else
return pb.and(
pb.equals("companyId", companyId),
pb.sql(codeSql.build())
);
}
ImdgPredicate finalPredicate = pb.and(
registryPredicate,
pb.notNull("object.tradingClearingRegistry"),
pb.and(
pb.greatEqual("eventTime", now),
pb.less("eventTime", now.plus(1, ChronoUnit.DAYS))
)
);
log.trace("sql predicate for registry_history {}", finalPredicate);
return finalPredicate;
}
@Override
public Optional<DepoBalanceRegister> isDuplicateInMap(Registry entity) {
return Optional.empty();
// Map<String, ? extends Comparable<?>> fieldValues = Map.of("companyId", entity.getCompanyId(), "sessionId", entity.getSessionId());
// Optional<EnumMessage> message = fieldValuesValidator.apply(fieldValues).tillFirstError();
// if (message.isPresent()) {
// log.warn("Illegal value in map for request to imdg: {}", message.get());
// throw new IllegalArgumentException(MessageFormat.format("Illegal value in request to imdg: {0}", message.get()));
// }
// return Optional.ofNullable(depoBalanceRegisterMap.getFirstObjectByFieldValues(fieldValues));
}
}

View file

@ -1,5 +1,8 @@
package ru.spcex.clearing.registry.service;
import org.apache.commons.lang3.tuple.MutablePair;
import org.apache.commons.lang3.tuple.MutableTriple;
import org.apache.commons.lang3.tuple.Triple;
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.producer.Producer;
import org.slf4j.Logger;
@ -7,49 +10,58 @@ import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.InitializingBean;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.account.Account;
import ru.clearing.classes.statics.data.company.Company;
import ru.clearing.classes.statics.data.company.CompanySymbols;
import ru.clearing.classes.statics.data.register.MoneyBalanceRegister;
import ru.clearing.classes.statics.data.registry.Registry;
import ru.clearing.classes.statics.data.registry.RegistryHistory;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
import ru.spcex.clearing.registry.util.PreClearMap;
import ru.spcex.platform.enumeration.AccountType;
import ru.spcex.platform.enumeration.CompanySymbol;
import ru.spcex.platform.enumeration.RegistryTradingParams;
import ru.spcex.platform.enumeration.RegistryUnit;
import ru.spcex.platform.imdg.api.Imdg;
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 ru.spcex.platform.utils.enumeration.IEnumKey;
import ru.spcex.platform.utils.validation.IValidator;
import java.math.BigDecimal;
import java.time.Instant;
import java.time.temporal.ChronoUnit;
import java.util.Collection;
import java.util.Map;
import java.util.*;
import java.util.function.Function;
import java.util.stream.Collectors;
import static ru.spcex.platform.enumeration.Task.createRegistry_GBRR;
@Service
public class MoneyBalanceRegisterService extends QueueConsumer implements InitializingBean {
public class MoneyBalanceRegisterService extends QueueConsumer implements InitializingBean, ICheckDuplicate<Registry, MoneyBalanceRegister> {
private final Logger log = LoggerFactory.getLogger(getClass());
private final Imdg<MoneyBalanceRegister> moneyBalanceRegisterMap;
private final Imdg<RegistryHistory> registryHistoryMap;
private final Imdg<Registry> registryMap;
private final Imdg<CompanySymbols> companySymbolsMap;
private final Imdg<Company> companyMap;
private final Imdg<Account> accountMap;
private final Function<Map<String, ?>, IValidator> fieldValuesValidator;
private final PreClearMap<MoneyBalanceRegister> preClearMap;
@Autowired
public MoneyBalanceRegisterService(Consumer<String, Object> kafkaQueue,
Producer<String, Object> kafkaProducer,
ImdgProvider imdgProvider) {
ImdgProvider imdgProvider,
Function<Map<String, ?>, IValidator> fieldValuesValidator) {
super(kafkaQueue, kafkaProducer);
this.moneyBalanceRegisterMap = imdgProvider.getImdg(IMDGDistributedNames.Map_MoneyBalanceRegister, MoneyBalanceRegister.class);
this.registryHistoryMap = imdgProvider.getImdg(IMDGDistributedNames.Map_RegistryHistory, RegistryHistory.class);
this.registryMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
this.companySymbolsMap = imdgProvider.getImdg(IMDGDistributedNames.Map_CompanySymbols, CompanySymbols.class);
this.companyMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Company, Company.class);
this.accountMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Account, Account.class);
this.fieldValuesValidator = fieldValuesValidator;
preClearMap = PreClearMap.instanceForInstantField(moneyBalanceRegisterMap, "created");
}
@ -61,70 +73,154 @@ public class MoneyBalanceRegisterService extends QueueConsumer implements Initia
init();
}
protected List<List<Registry>> groupingRegistry(Collection<Registry> registries) {
Map<MutablePair<Long, Long>, List<Registry>> regGroups = registries.stream().collect(Collectors.groupingBy(
(Registry reg) -> new MutablePair(reg.getCompanyId(), reg.getAccountId())
));
List<List<Registry>> result = new ArrayList<>(regGroups.values());
if (result.size() > 3) {
log.warn("Found too many ({}) registries in group: {}", result.size(), result);
}
return result;
// return regGroups.values().stream().map((List<Registry> regLst) -> {
// for (Registry r : regLst)
// if (RegistryUnit.T.equalsByKey(r.getRegistryUnit())) {
// return r;// приоритетно нужен AM_T
// }
// return regLst.iterator().next();
// }).collect(Collectors.toList());
}
public void moneyBalanceRegisterNew(BaseRequest<LauncherCommandRequest> userRequest) {
log.debug("LauncherCommandRequest received from {}", createRegistry_GBRR.topic());
ImdgPredicate predicate = getSqlForRegistryHistory(userRequest.getRequestPayload().getCompanyId());
Collection<RegistryHistory> registryHistories = registryHistoryMap.getCollectionObjectsByPredicate(predicate);
if (registryHistories.isEmpty()) {
log.debug("No registry history with such conditions {}", predicate);
ImdgPredicate sqlConditionForRegistry = getSqlForRegistries(userRequest.getRequestPayload().getCompanyId());
Collection<Registry> registries = registryMap.getCollectionObjectsByPredicate(sqlConditionForRegistry);
if (registries.isEmpty()) {
log.debug("No registry with such conditions {}", sqlConditionForRegistry);
return;
}
log.debug("Select {} Registry history by query {}", registryHistories.size(), predicate);
log.debug("Select {} Registry by query {}", registries.size(), sqlConditionForRegistry);
List<List<Registry>> registerGroups = groupingRegistry(registries);
preClearMap.preClearMap();
for (RegistryHistory registryHistory : registryHistories) {
if (BigDecimal.ZERO.compareTo(registryHistory.getObject().getDiffBalance()) != 0) continue;
insertMoneyBalanceRegister(registryHistory);
registerGroups.forEach((registryItems -> {
// if (isDuplicateInMap(registry).isEmpty()) {
insertMoneyBalanceRegister(registryItems);
// }
}
));
log.debug("successfully processed");
}
private void insertMoneyBalanceRegister(RegistryHistory registryHistory) {
private void insertMoneyBalanceRegister(List<Registry> registryGroup) {
log.trace("Started generating MoneyBalanceRegister entity...");
Registry registry = registryHistory.getObject();
MoneyBalanceRegister moneyBalanceRegister = new MoneyBalanceRegister();
Registry registryForRemainderSum = registryGroup.stream()
.filter(r -> RegistryUnit.T.equalsByKey(r.getRegistryUnit()))
.findAny().orElse(null);
Registry registryForBlockedSum = registryGroup.stream()
.filter(r -> RegistryUnit.B.equalsByKey(r.getRegistryUnit()))
.findAny().orElse(null);
Registry registryForUnblockedSum = registryGroup.stream()
.filter(r -> RegistryUnit.F.equalsByKey(r.getRegistryUnit()))
.findAny().orElse(null);
Registry registry = registryForRemainderSum;
if (registry == null)
registry = registryGroup.iterator().next();
log.trace("Group of {} registries: T={}; B={}; F={}", registryGroup.size(), registryForRemainderSum, registryForBlockedSum, registryForUnblockedSum);
// Registry registryForRemainderSum = registryMap.getFirstObjectByPredicate(getSqlForRemainderSum(registry.getCompanyId(), registry.getAccount()));
// Registry registryForBlockedSum = registryMap.getFirstObjectByPredicate(getSqlForBlockedSum(registry.getCompanyId(), registry.getAccount()));
// Registry registryForUnblockedSum = registryMap.getFirstObjectByPredicate(getSqlForUnblockedSum(registry.getCompanyId(), registry.getAccount()));
Company company = companyMap.getFirstObjectByFieldValues(Map.of("id", 2L));
CompanySymbols companySymbol = companySymbolsMap.getFirstObjectByFieldValues(Map.of("companySymbol", CompanySymbol.INN.getKey(), "companyId", registry.getCompanyId()));
String setHouseName = company != null ? company.getShortName() : null;
AccountType accountType = IEnumKey.getEnumByKey(AccountType.class, registry.getAccountType());
String inn = companySymbol != null ? companySymbol.getCompanySymbolValue() : null;
String accountType = registry.getAccountType();
String account = null;
if (accountType == AccountType.Clrn || accountType == AccountType.Info) {
String infoAccount = null;
if ("CLRN".equalsIgnoreCase(accountType)) {
Account accountValue = accountMap.getFirstObjectByFieldValues(Map.of("companyId", registry.getCompanyId(), "accountType", AccountType.Info.getKey()));
account = registry.getAccount();
infoAccount = accountValue != null ? accountValue.getAccount() : null;
} else if ("INFO".equalsIgnoreCase(accountType)) {
Account accountValue = accountMap.getFirstObjectByFieldValues(Map.of("companyId", registry.getCompanyId(), "accountType", AccountType.Clrn.getKey()));
account = accountValue != null ? accountValue.getAccount() : null;
infoAccount = registry.getAccount();
}
moneyBalanceRegister.setSetHouseName(setHouseName);
moneyBalanceRegister.setCreated(Instant.now());
moneyBalanceRegister.setCompanyId(registry.getCompanyId());
moneyBalanceRegister.setCompanyFullName(registry.getFullName());
moneyBalanceRegister.setAccount(account);
moneyBalanceRegister.setRemainderSum(registry.getBalance());
moneyBalanceRegister.setInfoAccount(infoAccount);
moneyBalanceRegister.setRemainderSum(registryForRemainderSum != null ? registryForRemainderSum.getBalance() : null);
moneyBalanceRegister.setBlockedSum(registryForBlockedSum != null ? registryForBlockedSum.getBalance() : null);
moneyBalanceRegister.setUnblockedSum(registryForUnblockedSum != null ? registryForUnblockedSum.getBalance() : null);
moneyBalanceRegister.setInn(inn);
moneyBalanceRegister.setSessionId(registry.getSessionId());
moneyBalanceRegister.setCompanyId(registry.getCompanyId());
moneyBalanceRegister.setCreated(Instant.now());
moneyBalanceRegister.setUpdated(moneyBalanceRegister.getCreated());
moneyBalanceRegisterMap.insert(moneyBalanceRegister);
log.debug("inserted successfully MoneyBalanceRegister entity with id: {}", moneyBalanceRegister.getId());
}
protected ImdgPredicate getSqlForRegistryHistory(Long companyId) {
Instant now = Instant.now().truncatedTo(ChronoUnit.DAYS);
ImdgPredicateBuilder pb = registryHistoryMap.predicateBuilder();
RegistryCodeSqlBuilder codeSql = RegistryCodeSqlBuilder.getInstance(RegistryTradingParams.AM_T).prefix("object");
ImdgPredicate registryPredicate;
protected ImdgPredicate getSqlForRegistries(Long companyId) {
ImdgPredicateBuilder pb = registryMap.predicateBuilder();
ImdgPredicate prdct;
if (companyId == null) {
registryPredicate = pb.sql(codeSql.build());
prdct = pb.sql(RegistryCodeSqlBuilder.getInstance(RegistryTradingParams.AM__).build());
} else {
registryPredicate = pb.and(
pb.equals("object.companyId", companyId),
pb.sql(codeSql.build())
prdct = pb.and(
pb.equals("companyId", companyId),
pb.sql(RegistryCodeSqlBuilder.getInstance(RegistryTradingParams.AM__).build())
);
}
ImdgPredicate finalPredicate = pb.and(
registryPredicate,
pb.notNull("object.tradingClearingRegistry"),
pb.and(
pb.greatEqual("eventTime", now),
pb.less("eventTime", now.plus(1, ChronoUnit.DAYS))
)
);
log.trace("sql predicate for registry_history {}", finalPredicate);
return finalPredicate;
log.trace("sql predicate for registries {}", prdct);
return prdct;
}
protected ImdgPredicate getSqlForRemainderSum(Long companyId, String account) {
ImdgPredicateBuilder pb = registryMap.predicateBuilder();
ImdgPredicate prdct = pb.and(
pb.equals("companyId", companyId),
pb.equals("account", account),
pb.sql(RegistryCodeSqlBuilder.getInstance(RegistryTradingParams.AM_T).build())
);
log.trace("sql predicate for registries remainder sum {}", prdct);
return prdct;
}
protected ImdgPredicate getSqlForBlockedSum(Long companyId, String account) {
ImdgPredicateBuilder pb = registryMap.predicateBuilder();
ImdgPredicate prdct = pb.and(
pb.equals("companyId", companyId),
pb.equals("account", account),
pb.sql(RegistryCodeSqlBuilder.getInstance(RegistryTradingParams.AM_B).build())
);
log.trace("sql predicate for registries blocked sum {}", prdct);
return prdct;
}
protected ImdgPredicate getSqlForUnblockedSum(Long companyId, String account) {
ImdgPredicateBuilder pb = registryMap.predicateBuilder();
ImdgPredicate prdct = pb.and(
pb.equals("companyId", companyId),
pb.equals("account", account),
pb.sql(RegistryCodeSqlBuilder.getInstance(RegistryTradingParams.AM_F).build())
);
log.trace("sql predicate for registries unblocked sum {}", prdct);
return prdct;
}
@Override
public Optional<MoneyBalanceRegister> isDuplicateInMap(Registry entity) {
return Optional.empty();
// Map<String, Long> fieldValues = new HashMap<>();
// fieldValues.put("companyId", entity.getCompanyId());
// fieldValues.put("sessionId", entity.getSessionId());
// Optional<EnumMessage> message = fieldValuesValidator.apply(fieldValues).tillFirstError();
// if (message.isPresent()) {
// log.warn("Illegal value in map for request to imdg: {}", message.get());
// throw new IllegalArgumentException(MessageFormat.format("Illegal value in request to imdg: {0}", message.get()));
// }
// return Optional.ofNullable(moneyBalanceRegisterMap.getFirstObjectByFieldValues(fieldValues));
}
}

View file

@ -13,7 +13,6 @@ import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit.jupiter.SpringExtension;
import ru.clearing.classes.statics.data.register.DepoBalanceRegister;
import ru.clearing.classes.statics.data.registry.Registry;
import ru.clearing.classes.statics.data.registry.RegistryHistory;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
import ru.spcex.clearing.registry.config.ValidationConfig;
@ -26,7 +25,6 @@ import ru.spcex.platform.imdg.api.ImdgProvider;
import javax.annotation.PostConstruct;
import java.math.BigDecimal;
import java.time.Instant;
import java.util.Collection;
import static org.junit.jupiter.api.Assertions.assertEquals;
@ -50,7 +48,7 @@ class DepoBalanceRegisterServiceTest {
private static final String SECURITY_SYMBOL= "SECURITY_SYMBOL";
private Imdg<DepoBalanceRegister> depoBalanceRegisterMap;
private Imdg<RegistryHistory> registryHistoryMap;
private Imdg<Registry> registryMap;
@Autowired
private DepoBalanceRegisterService depoBalanceRegisterService;
@ -68,29 +66,27 @@ class DepoBalanceRegisterServiceTest {
ImdgTestConfig.waitAvailableImdgProviderAndAddAdminWithDefaultId();
//initializing of imdg
registryHistoryMap = imdgProvider.getImdg(IMDGDistributedNames.Map_RegistryHistory, RegistryHistory.class);
registryMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
depoBalanceRegisterMap = imdgProvider.getImdg(IMDGDistributedNames.Map_DepoBalanceRegister, DepoBalanceRegister.class);
// initialing objects for imdg repos
Long securityId = 1L;
Registry registry = new Registry();
registry.setRegistryDesignation(RegistryDesignation.A.getKey());
registry.setRegistryInstrumentType(RegistryInstrumentType.S.getKey());
registry.setRegistryUnit(RegistryUnit.T.getKey());
registry.setDiffBalance(BigDecimal.ZERO);
registry.setTradingClearingRegistry("not null value");
registry.setRegistryStatus(RegistryStatus.OK.getKey());
registry.setSessionId(sessionId);
registry.setCompanyId(companyId);
registry.setAccountType(AccountType.Depo.getKey());
registry.setAccount(account);
registry.setBalance(balance);
registry.setSecuritySymbol(SECURITY_SYMBOL);
RegistryHistory registryHistory = new RegistryHistory();
registryHistory.setObject(registry);
registryHistory.setEventTime(Instant.now());
registryHistoryMap.insert(registryHistory);
registryMap.insert(registry);
}
@Test
@Order(1)
void executionRegisterNewForEventTime() {
void executionRegisterNewForClearingDate() {
//making launcherCommand
LauncherCommandRequest launcherCommandRequest = new LauncherCommandRequest();
@ -108,6 +104,7 @@ class DepoBalanceRegisterServiceTest {
//take all values is procceed imdg
Collection<DepoBalanceRegister> allValues = depoBalanceRegisterMap.getAllValues();
DepoBalanceRegister depoBalanceRegister = allValues.iterator().next();
assertEquals(sessionId, depoBalanceRegister.getSessionId());
assertEquals(companyId, depoBalanceRegister.getCompanyId());
assertEquals(account, depoBalanceRegister.getDepoCode());
assertEquals(balance, depoBalanceRegister.getQuantity());

View file

@ -11,10 +11,11 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit.jupiter.SpringExtension;
import ru.clearing.classes.statics.data.account.Account;
import ru.clearing.classes.statics.data.company.Company;
import ru.clearing.classes.statics.data.company.CompanySymbols;
import ru.clearing.classes.statics.data.register.MoneyBalanceRegister;
import ru.clearing.classes.statics.data.registry.Registry;
import ru.clearing.classes.statics.data.registry.RegistryHistory;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
import ru.spcex.clearing.registry.config.ValidationConfig;
@ -27,7 +28,6 @@ import ru.spcex.platform.imdg.api.ImdgProvider;
import javax.annotation.PostConstruct;
import java.math.BigDecimal;
import java.time.Instant;
import java.time.LocalDate;
import java.util.Collection;
@ -62,8 +62,10 @@ class MoneyBalanceRegisterServiceTest {
@Autowired
private MoneyBalanceRegisterService moneyBalanceRegisterService;
private Imdg<MoneyBalanceRegister> moneyBalanceRegisterMap;
private Imdg<RegistryHistory> registryHistoryMap;
private Imdg<Registry> registryMap;
private Imdg<CompanySymbols> companySymbolsMap;
private Imdg<Company> companyMap;
private Imdg<Account> accountMap;
@Autowired
@Qualifier("hazelcastServiceTest")
private ImdgProvider imdgProvider;
@ -77,16 +79,16 @@ class MoneyBalanceRegisterServiceTest {
ImdgTestConfig.waitAvailableImdgProviderAndAddAdminWithDefaultId();
moneyBalanceRegisterMap = imdgProvider.getImdg(IMDGDistributedNames.Map_MoneyBalanceRegister, MoneyBalanceRegister.class);
registryMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
companySymbolsMap = imdgProvider.getImdg(IMDGDistributedNames.Map_CompanySymbols, CompanySymbols.class);
companyMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Company, Company.class);
registryHistoryMap = imdgProvider.getImdg(IMDGDistributedNames.Map_RegistryHistory, RegistryHistory.class);
accountMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Account, Account.class);
}
@Test
@Order(1)
void executionRegisterNewForClearingDate() {
Registry registry = new Registry();
registry.setTradingClearingRegistry("not null value");
registry.setDiffBalance(BigDecimal.ZERO);
registry.setRegistryUnit(RegistryUnit.T.getKey());
registry.setRegistryInstrumentType(REGISTRY_INSTRUMENT_TYPE);
registry.setRegistryDesignation(REGISTRY_DESIGNATION);
@ -100,16 +102,29 @@ class MoneyBalanceRegisterServiceTest {
registry.setBalance(REGISTRY_BALANCE);
registry.setSettlementDate(REGISTRY_SETTLEMENT_DATE);
registry.setSecurityId(REGISTRY_SECURITY_ID);
RegistryHistory registryHistory = new RegistryHistory();
registryHistory.setObject(registry);
registryHistory.setEventTime(Instant.now());
registryHistoryMap.insert(registryHistory);
registryMap.insert(registry);
registry.setId(null);
registry.setRegistryUnit(RegistryUnit.B.getKey());
registryMap.insert(registry);
registry.setId(null);
registry.setRegistryUnit(RegistryUnit.F.getKey());
registryMap.insert(registry);
Company company = new Company();
company.setId(2L);
company.setShortName(COMPANY_SHORTNAME);
companyMap.insert(company);
CompanySymbols companySymbols = new CompanySymbols();
companySymbols.setCompanyId(REGISTRY_COMPANY_ID);
companySymbols.setCompanySymbol(COMPANY_SYMBOL_INN);
companySymbols.setCompanySymbolValue(COMPANY_SYMBOL_VALUE);
companySymbolsMap.insert(companySymbols);
Account account = new Account();
account.setCompanyId(REGISTRY_COMPANY_ID);
account.setAccountType(AccountType.Info.getKey());
account.setAccount(ACCOUNT_ACCOUNT);
accountMap.insert(account);
//making launcherCommand
LauncherCommandRequest launcherCommandRequest = new LauncherCommandRequest();
launcherCommandRequest.setCompanyId(REGISTRY_COMPANY_ID);
@ -127,8 +142,13 @@ class MoneyBalanceRegisterServiceTest {
Collection<MoneyBalanceRegister> allValues = moneyBalanceRegisterMap.getAllValues();
MoneyBalanceRegister moneyBalanceRegister = allValues.iterator().next();
assertEquals(REGISTRY_COMPANY_ID, moneyBalanceRegister.getCompanyId());
assertEquals(REGISTRY_FULL_NAME, moneyBalanceRegister.getCompanyFullName());
assertEquals(REGISTRY_BALANCE, moneyBalanceRegister.getRemainderSum());
assertEquals(REGISTRY_BALANCE, moneyBalanceRegister.getBlockedSum());
assertEquals(REGISTRY_BALANCE, moneyBalanceRegister.getUnblockedSum());
assertEquals(REGISTRY_ACCOUNT, moneyBalanceRegister.getAccount());
assertEquals(ACCOUNT_ACCOUNT, moneyBalanceRegister.getInfoAccount());
assertEquals(COMPANY_SHORTNAME, moneyBalanceRegister.getSetHouseName());
assertEquals(COMPANY_SYMBOL_VALUE, moneyBalanceRegister.getInn());
}
}

View file

@ -7,7 +7,6 @@ import com.opencsv.bean.CsvNumber;
import java.math.BigDecimal;
import java.time.LocalDate;
import java.time.LocalDateTime;
public class KSCommissionTradesReport {
@CsvBindByName(column = "FIRM_ID")
@ -74,19 +73,6 @@ public class KSCommissionTradesReport {
@CsvBindByPosition(position = 14)
private Long sessionId;
@CsvBindByName(column = "SESSION_DATE")
@CsvBindByPosition(position = 15)
@CsvDate(value = "dd.MM.yyyy'T'HH:mm:ss")
private LocalDateTime sessionDate;
@CsvBindByName(column = "TKR_INITIATOR")
@CsvBindByPosition(position = 16)
private String tkrInitiator;
@CsvBindByName(column = "ACCOUNT_INITIATOR")
@CsvBindByPosition(position = 17)
private String accountInitiator;
public String getFirmId() {
return firmId;
@ -207,28 +193,4 @@ public class KSCommissionTradesReport {
public void setSessionId(Long sessionId) {
this.sessionId = sessionId;
}
public LocalDateTime getSessionDate() {
return sessionDate;
}
public void setSessionDate(LocalDateTime sessionDate) {
this.sessionDate = sessionDate;
}
public String getTkrInitiator() {
return tkrInitiator;
}
public void setTkrInitiator(String tkrInitiator) {
this.tkrInitiator = tkrInitiator;
}
public String getAccountInitiator() {
return accountInitiator;
}
public void setAccountInitiator(String accountInitiator) {
this.accountInitiator = accountInitiator;
}
}

View file

@ -7,7 +7,6 @@ import com.opencsv.bean.CsvDate;
import com.opencsv.bean.CsvNumber;
import java.math.BigDecimal;
import java.time.LocalDate;
import java.time.LocalDateTime;
public class KSRepTradesReport {
@ -71,11 +70,6 @@ public class KSRepTradesReport {
@CsvBindByPosition(position = 13)
private String status;
@CsvBindByName(column = "REPAYM_DATE")
@CsvBindByPosition(position = 14)
@CsvDate(value = "dd.MM.yyyy")
private LocalDate repaymDate;
public Long getSessionId() {
return sessionId;
@ -188,12 +182,4 @@ public class KSRepTradesReport {
public void setDepoAccount(String depoAccount) {
this.depoAccount = depoAccount;
}
public LocalDate getRepaymDate() {
return repaymDate;
}
public void setRepaymDate(LocalDate repaymDate) {
this.repaymDate = repaymDate;
}
}

View file

@ -4,7 +4,6 @@ import org.springframework.stereotype.Component;
import ru.clearing.classes.statics.data.account.Account;
import ru.clearing.classes.statics.data.company.Company;
import ru.clearing.classes.statics.data.execution.ExecutionDeposit;
import ru.clearing.classes.statics.data.misc.Session;
import ru.clearing.classes.statics.data.registry.Registry;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.reports.builders.CSVReportBuilder;
@ -19,7 +18,6 @@ import ru.spcex.platform.imdg.api.predicate.specific.RegistryCodeSqlBuilder;
import ru.spcex.platform.utils.enumeration.IEnumKey;
import java.time.LocalDate;
import java.time.LocalDateTime;
import java.time.format.DateTimeFormatter;
import java.util.ArrayList;
import java.util.Collection;
@ -34,7 +32,6 @@ public class KSCommissionTradesReportBuilder extends CSVReportBuilder<EmptyParam
private final Imdg<Company> companyImdg;
private final Imdg<ExecutionDeposit> executionDepositImdg;
private final Imdg<Account> accountImdg;
private final Imdg<Session> sessionImdg;
private List<KSCommissionTradesReport> rows = null;
@ -43,7 +40,6 @@ public class KSCommissionTradesReportBuilder extends CSVReportBuilder<EmptyParam
companyImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Company, Company.class);
executionDepositImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionDeposit, ExecutionDeposit.class);
accountImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Account, Account.class);
sessionImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Session, Session.class);
}
@Override
@ -177,39 +173,6 @@ public class KSCommissionTradesReportBuilder extends CSVReportBuilder<EmptyParam
account = accountFromImdg.getAccount();
}
String tkrInitiator = null;
String accountInitiator = null;
ImdgPredicate tmtPredicate = RegistryCodeSqlBuilder.getInstance(RegistryTradingParams.TM_T).buildPredicate(pb);
Registry tmtRegistry = registryImdg.getFirstObjectByPredicate(
pb.and(
pb.equals("counterPartyId", registry.getCompanyId()),
pb.equals("groupId", registry.getGroupId()),
tmtPredicate
)
);
if (tmtRegistry != null) {
tkrInitiator = tmtRegistry.getTradingClearingRegistry();
accountType = IEnumKey.getEnumByKey(AccountType.class, tmtRegistry.getAccountType());
if (accountType == AccountType.Clrn) {
accountInitiator = tmtRegistry.getAccount();
} else if (accountType == AccountType.Info) {
Account accountFromImdg = accountImdg.getFirstObjectByFieldValues(Map.of(
"accountType", AccountType.Anlt.getKey(),
"companyId", 1L
));
accountInitiator = accountFromImdg.getAccount();
}
}
Session session = sessionImdg.getSingleObjectByID(registry.getSessionId());
LocalDateTime sessionDate = null;
Long sessionId = null;
if (registryStatus != RegistryStatus.MNG && registryStatus != RegistryStatus.PROC) {
sessionId = registry.getSessionId();
if (session != null) {
sessionDate = LocalDateTime.ofInstant(session.getCreated(), getReportZoneId());
}
}
KSCommissionTradesReport ksCommissionTradesReport = new KSCommissionTradesReport();
ksCommissionTradesReport.setFirmId(firmId);
ksCommissionTradesReport.setFirmIdInitiator(firmIdInitiator);
@ -228,10 +191,7 @@ public class KSCommissionTradesReportBuilder extends CSVReportBuilder<EmptyParam
ksCommissionTradesReport.setTrId(registry.getGroupId());
ksCommissionTradesReport.setTkr(registry.getTradingClearingRegistry());
ksCommissionTradesReport.setAccount(account);
ksCommissionTradesReport.setSessionId(sessionId);
ksCommissionTradesReport.setSessionDate(sessionDate);
ksCommissionTradesReport.setTkrInitiator(tkrInitiator);
ksCommissionTradesReport.setAccountInitiator(accountInitiator);
ksCommissionTradesReport.setSessionId(registry.getSessionId());
rows.add(ksCommissionTradesReport);
} catch (Throwable e) {

View file

@ -22,7 +22,10 @@ import java.math.BigDecimal;
import java.time.Instant;
import java.time.LocalDate;
import java.time.temporal.ChronoUnit;
import java.util.*;
import java.util.ArrayList;
import java.util.Collection;
import java.util.List;
import java.util.Map;
@Component
public class KSRepCashRegistersReportBuilder extends CSVReportBuilder<EmptyParams, KSRepCashRegistersReport> {
@ -150,27 +153,9 @@ public class KSRepCashRegistersReportBuilder extends CSVReportBuilder<EmptyParam
pb.equals("registryStatus", RegistryStatus.CLRD.getKey())
)
);
Set<String> contracts = new HashSet<>();
BigDecimal cmtBalanceSum = BigDecimal.ZERO;
for (Registry cmtRegistry : cmtRegistries) {
BigDecimal balance = cmtRegistry.getBalance();
cmtBalanceSum = cmtBalanceSum.add(balance);
contracts.add(cmtRegistry.getContract());
}
ImdgPredicate dmxPredicate = RegistryCodeSqlBuilder.getInstance(RegistryTradingParams.DM_X).buildPredicate(pb);
Collection<Registry> dmxRegistries = registryImdg.getCollectionObjectsByPredicate(
pb.and(
dmxPredicate,
pb.in("contract", contracts.toArray(new String[0])),
pb.equals("tradingCode", registry.getTradingCode()),
pb.equals("companyId", registry.getCompanyId()),
pb.equals("account", registry.getAccount()),
pb.equals("registryStatus", RegistryStatus.OK.getKey())
)
);
BigDecimal dmxBalanceSum = dmxRegistries.stream().map(Registry::getBalance).reduce(BigDecimal.ZERO, BigDecimal::add);
BigDecimal cmtBalanceSum = cmtRegistries.stream().map(Registry::getBalance).reduce(BigDecimal.ZERO, BigDecimal::add);
BigDecimal lmtBalanceSum = lmtRegistries.stream().map(Registry::getBalance).reduce(BigDecimal.ZERO, BigDecimal::add);
ksRepCashRegistersReport.setChangeBalance(cmtBalanceSum.subtract(lmtBalanceSum).subtract(dmxBalanceSum));
ksRepCashRegistersReport.setChangeBalance(cmtBalanceSum.subtract(lmtBalanceSum));
registryPredicate = RegistryCodeSqlBuilder.getInstance(RegistryTradingParams.DM_I).buildPredicate(pb);

View file

@ -209,7 +209,6 @@ public class KSRepTradesReportBuilder extends CSVReportBuilder<EmptyParams, KSRe
}
ksRepTradesReport.setCashLiability(firstLegAmount);
ksRepTradesReport.setDepoLiability(BigDecimal.ZERO);
ksRepTradesReport.setRepaymDate(executionDeposit.getSecondLegSettlementDate());
} else if (execution instanceof ExecutionFond executionFond) {
if (tkr.getDepoAccountId() != null) {

View file

@ -192,7 +192,6 @@ public class ReportServiceTest_KS {
cmtRegistry_1.setAccount("ACCOUNT");
cmtRegistry_1.setRegistryStatus(RegistryStatus.CLRD.getKey());
cmtRegistry_1.setBalance(BigDecimal.valueOf(111L));
cmtRegistry_1.setContract("CONTRACT_1");
registryImdg.insert(cmtRegistry_1);
Registry cmtRegistry_2 = new Registry();
@ -206,7 +205,6 @@ public class ReportServiceTest_KS {
cmtRegistry_2.setAccount("ACCOUNT");
cmtRegistry_2.setRegistryStatus(RegistryStatus.CLRD.getKey());
cmtRegistry_2.setBalance(BigDecimal.valueOf(111L));
cmtRegistry_2.setContract("CONTRACT_2");
registryImdg.insert(cmtRegistry_2);
Registry cmtRegistry_3 = new Registry();
@ -261,34 +259,6 @@ public class ReportServiceTest_KS {
lmtRegistry_3.setBalance(BigDecimal.valueOf(11L));
registryImdg.insert(lmtRegistry_3);
Registry dmxRegistry_1 = new Registry();
dmxRegistry_1.setRegistryCode("DM_X");
dmxRegistry_1.setRegistryDesignation(RegistryDesignation.D.getKey());
dmxRegistry_1.setRegistryInstrumentType(RegistryInstrumentType.M.getKey());
dmxRegistry_1.setRegistryUnit(RegistryUnit.X.getKey());
dmxRegistry_1.setSessionId(SESSION_ID_FOR_DAY);
dmxRegistry_1.setTradingCode("TRADING_CODE");
dmxRegistry_1.setCompanyId(COMPANY_ID);
dmxRegistry_1.setAccount("ACCOUNT");
dmxRegistry_1.setContract("CONTRACT_1");
dmxRegistry_1.setRegistryStatus(RegistryStatus.OK.getKey());
dmxRegistry_1.setBalance(BigDecimal.valueOf(11L));
registryImdg.insert(dmxRegistry_1);
Registry dmxRegistry_2 = new Registry();
dmxRegistry_2.setRegistryCode("DM_X");
dmxRegistry_2.setRegistryDesignation(RegistryDesignation.D.getKey());
dmxRegistry_2.setRegistryInstrumentType(RegistryInstrumentType.M.getKey());
dmxRegistry_2.setRegistryUnit(RegistryUnit.X.getKey());
dmxRegistry_2.setSessionId(SESSION_ID_FOR_DAY);
dmxRegistry_2.setTradingCode("TRADING_CODE");
dmxRegistry_2.setCompanyId(COMPANY_ID);
dmxRegistry_2.setAccount("ACCOUNT");
dmxRegistry_2.setContract("CONTRACT_2");
dmxRegistry_2.setRegistryStatus(RegistryStatus.OK.getKey());
dmxRegistry_2.setBalance(BigDecimal.valueOf(11L));
registryImdg.insert(dmxRegistry_2);
Registry addWithdrawRegistry_1 = new Registry();
addWithdrawRegistry_1.setCompanyId(COMPANY_ID);
@ -355,8 +325,6 @@ public class ReportServiceTest_KS {
registryImdg.delete(lmtRegistry_1);
registryImdg.delete(lmtRegistry_2);
registryImdg.delete(lmtRegistry_3);
registryImdg.delete(dmxRegistry_1);
registryImdg.delete(dmxRegistry_2);
sessionImdg.delete(sessionForDay);
accountImdg.delete(account);
}
@ -729,7 +697,6 @@ public class ReportServiceTest_KS {
executionDeposit.setExchangeExecutionId(222L);
executionDeposit.setExchangeExecutionTime(Instant.ofEpochMilli(LocalDateTime.of(2022, 1, 1, 1, 1, 1).toEpochSecond(ZoneOffset.UTC)));
executionDeposit.setCoverageStatus(CoverageStatus.DEND.getKey());
executionDeposit.setSecondLegSettlementDate(LocalDate.of(2000, 1, 1));
executionDepositImdg.insert(executionDeposit);
String jsonString = TestUtils.getJsonStringForSystem(reportRequest, 0L);
@ -1064,10 +1031,6 @@ public class ReportServiceTest_KS {
companySymbols.setCompanySymbolValue("COMPANY_SYMBOL_VALUE");
companySymbolsImdg.insert(companySymbols);
Session session = new Session();
session.setCreated(LocalDateTime.of(2000, 1, 1, 0, 0).toInstant(ZoneOffset.UTC));
Long sessionId = sessionImdg.insert(session);
Registry registry = new Registry();
registry.setRegistryStatus(RegistryStatus.CLRD.getKey());
registry.setCompanyId(COMPANY_ID);
@ -1090,7 +1053,7 @@ public class ReportServiceTest_KS {
registry.setTradingClearingRegistry("TRADING_CLEARING_REGISTRY");
registry.setAccountType(AccountType.Clrn.getKey());
registry.setAccount("ACCOUNT");
registry.setSessionId(sessionId);
registry.setSessionId(1L);
registryImdg.insert(registry);
Registry omtRegistry = new Registry();
@ -1102,18 +1065,6 @@ public class ReportServiceTest_KS {
omtRegistry.setGroupId(23071812300L);
registryImdg.insert(omtRegistry);
Registry tmtRegistry = new Registry();
tmtRegistry.setRegistryCode("TM_T");
tmtRegistry.setAccount("TMT_ACCOUNT");
tmtRegistry.setAccountType(AccountType.Clrn.getKey());
tmtRegistry.setRegistryDesignation(RegistryDesignation.T.getKey());
tmtRegistry.setRegistryInstrumentType(RegistryInstrumentType.M.getKey());
tmtRegistry.setRegistryUnit(RegistryUnit.T.getKey());
tmtRegistry.setCounterPartyId(COMPANY_ID);
tmtRegistry.setTradingClearingRegistry("TKR_TMT");
tmtRegistry.setGroupId(23071812300L);
registryImdg.insert(tmtRegistry);
ExecutionDeposit executionDeposit = new ExecutionDeposit();
executionDeposit.setExchangeExecutionId(registry.getGroupId());
executionDeposit.setMarket("market");
@ -1131,12 +1082,10 @@ public class ReportServiceTest_KS {
File outFile = outFolder.listFiles()[0];
File expectedFile = new File(getClass().getClassLoader().getResource("expected_reports_csv/ks/KS_COMMISSION_TRADES_expected.csv").getFile());
compareCSVFiles(expectedFile, outFile, Arrays.asList("AGREEMENT_SETTLEMENT_DATE", "REPAYM_DATE", "SESSION_ID"));
compareCSVFiles(expectedFile, outFile, Arrays.asList("AGREEMENT_SETTLEMENT_DATE", "REPAYM_DATE"));
registryImdg.delete(registry);
registryImdg.delete(omtRegistry);
registryImdg.delete(tmtRegistry);
sessionImdg.delete(session);
companySymbolsImdg.delete(companySymbols);
executionDepositImdg.delete(executionDeposit);
}

View file

@ -1,2 +1,2 @@
"FIRM_ID","FIRM_ID_INITIATOR","SECCODE","CLASS_CODE","AGREEMENT_NUMBER","AGREEMENT_NUMBER_POSTFIX","AGREEMENT_DATE","AGREEMENT_SETTLEMENT_DATE","SUM","REPAYM_DATE","EXEC_STATUS","TR_ID","TKR","ACCOUNT","SESSION_ID","SESSION_DATE","TKR_INITIATOR","ACCOUNT_INITIATOR"
"CLEARING_CODE","TRADING_CODE","CONT~RACT","market","CONT","CONT~RACT","01.01.1991","28.11.2023","-77.78","08.12.2023","Исполнен","23071812300","TRADING_CLEARING_REGISTRY","ACCOUNT","3","01.01.2000T03:00:00","TKR_TMT","TMT_ACCOUNT"
"FIRM_ID","FIRM_ID_INITIATOR","SECCODE","CLASS_CODE","AGREEMENT_NUMBER","AGREEMENT_NUMBER_POSTFIX","AGREEMENT_DATE","AGREEMENT_SETTLEMENT_DATE","SUM","REPAYM_DATE","EXEC_STATUS","TR_ID","TKR","ACCOUNT","SESSION_ID"
"CLEARING_CODE","TRADING_CODE","CONT~RACT","market","CONT","CONT~RACT","01.01.1991","14.11.2023","-77.78","24.11.2023","Исполнен","23071812300","TRADING_CLEARING_REGISTRY","ACCOUNT","1"

1 FIRM_ID FIRM_ID_INITIATOR SECCODE CLASS_CODE AGREEMENT_NUMBER AGREEMENT_NUMBER_POSTFIX AGREEMENT_DATE AGREEMENT_SETTLEMENT_DATE SUM REPAYM_DATE EXEC_STATUS TR_ID TKR ACCOUNT SESSION_ID SESSION_DATE TKR_INITIATOR ACCOUNT_INITIATOR
2 CLEARING_CODE TRADING_CODE CONT~RACT market CONT CONT~RACT 01.01.1991 28.11.2023 14.11.2023 -77.78 08.12.2023 24.11.2023 Исполнен 23071812300 TRADING_CLEARING_REGISTRY ACCOUNT 3 1 01.01.2000T03:00:00 TKR_TMT TMT_ACCOUNT

View file

@ -1,3 +1,3 @@
"FIRMID","ACCOUNT","TKR","DATE","REGISTER_CODE","REGISTER_NAME","OPEN_BALANCE","CHANGE_BALANCE","ADD_WITHDRAW","CLOSE_BALANCE","REMARKS","TRADE_NUM"
"TRADING_CODE","ACCOUNT_FROM_ACCOUNT","TRADING_CLEARING_REGISTRY","01.11.2023","AM_T","AM_T_NAME","77.78",278.00,"111.00","466.78","",""
"TRADING_CODE","ACCOUNT","TRADING_CLEARING_REGISTRY","01.11.2023","AM_T","AM_T_NAME","77.78",278.00,"111.00",466.78,"",""
"TRADING_CODE","ACCOUNT_FROM_ACCOUNT","TRADING_CLEARING_REGISTRY","01.11.2023","AM_T","AM_T_NAME","77.78","300.00","111.00","488.78","",""
"TRADING_CODE","ACCOUNT","TRADING_CLEARING_REGISTRY","01.11.2023","AM_T","AM_T_NAME","77.78","300.00","111.00","488.78","",""

1 FIRMID ACCOUNT TKR DATE REGISTER_CODE REGISTER_NAME OPEN_BALANCE CHANGE_BALANCE ADD_WITHDRAW CLOSE_BALANCE REMARKS TRADE_NUM
2 TRADING_CODE ACCOUNT_FROM_ACCOUNT TRADING_CLEARING_REGISTRY 01.11.2023 AM_T AM_T_NAME 77.78 278.00 300.00 111.00 466.78 488.78
3 TRADING_CODE ACCOUNT TRADING_CLEARING_REGISTRY 01.11.2023 AM_T AM_T_NAME 77.78 278.00 300.00 111.00 466.78 488.78

View file

@ -1,5 +1,5 @@
"SESSION_ID","SESSION_DATE","TR_ID","AGREEMENT_NUMBER","FIRM_ID","TKR","ACCOUNT","DEPO_ACCOUNT","CASH_LIABILITY","DEPO_LIABILITY","TRADE_NUM","CLASS_CODE","TRADE_DATE","STATUS","REPAYM_DATE"
"7","20.01.1970T07:09:07","222","CONTRACT","TRADING_CODE","TKR",,,"8888.89","0","222","3333","20.01.1970T02:49:58","2","01.01.2000"
"6","20.01.1970T15:54:43","222",,"TRADING_CODE","TKR",,"ACCOUNT","-8888.89","123","222","3333","20.01.1970T02:49:58","1",
"6","20.01.1970T15:54:43","222",,"TRADING_CODE","TKR",,"ACC_SYMBOL_VALUE","8888.89","-123","222","3333","20.01.1970T02:49:58","2",
"6","20.01.1970T15:54:43","222",,"TRADING_CODE","TKR",,"ACC_SYMBOL_VALUE","-8888.89","123","222","3333","20.01.1970T02:49:58","2",
"SESSION_ID","SESSION_DATE","TR_ID","AGREEMENT_NUMBER","FIRM_ID","TKR","ACCOUNT","DEPO_ACCOUNT","CASH_LIABILITY","DEPO_LIABILITY","TRADE_NUM","CLASS_CODE","TRADE_DATE","STATUS"
"7","20.01.1970T07:09:07","222","CONTRACT","TRADING_CODE","TKR","","","8888.89","0","222","3333","20.01.1970T02:49:58","2"
"6","20.01.1970T15:54:43","222","","TRADING_CODE","TKR","","ACCOUNT","-8888.89","123","222","3333","20.01.1970T02:49:58","1"
"6","20.01.1970T15:54:43","222","","TRADING_CODE","TKR","","ACC_SYMBOL_VALUE","8888.89","-123","222","3333","20.01.1970T02:49:58","2"
"6","20.01.1970T15:54:43","222","","TRADING_CODE","TKR","","ACC_SYMBOL_VALUE","-8888.89","123","222","3333","20.01.1970T02:49:58","2"

1 SESSION_ID SESSION_DATE TR_ID AGREEMENT_NUMBER FIRM_ID TKR ACCOUNT DEPO_ACCOUNT CASH_LIABILITY DEPO_LIABILITY TRADE_NUM CLASS_CODE TRADE_DATE STATUS REPAYM_DATE
2 7 20.01.1970T07:09:07 222 CONTRACT TRADING_CODE TKR 8888.89 0 222 3333 20.01.1970T02:49:58 2 01.01.2000
3 6 20.01.1970T15:54:43 222 TRADING_CODE TKR ACCOUNT -8888.89 123 222 3333 20.01.1970T02:49:58 1
4 6 20.01.1970T15:54:43 222 TRADING_CODE TKR ACC_SYMBOL_VALUE 8888.89 -123 222 3333 20.01.1970T02:49:58 2
5 6 20.01.1970T15:54:43 222 TRADING_CODE TKR ACC_SYMBOL_VALUE -8888.89 123 222 3333 20.01.1970T02:49:58 2

View file

@ -1,18 +0,0 @@
package ru.spcex.platform.enumeration;
import ru.spcex.platform.utils.enumeration.IEnumKey;
public enum ClearingAccountType implements IEnumKey {
A("A"), B("B");
private final String key;
ClearingAccountType(String key) {
this.key = key;
}
@Override
public String getKey() {
return key;
}
}

View file

@ -12,7 +12,6 @@ import java.util.stream.Collectors;
public class RegistryCodeSqlBuilder {
private List<RegistryTradingParams> registryTradingParams;
private String prefix = null;
public static RegistryCodeSqlBuilder getInstance(RegistryTradingParams... tradingParams) {
return new RegistryCodeSqlBuilder(Arrays.asList(tradingParams));
@ -24,21 +23,11 @@ public class RegistryCodeSqlBuilder {
public String build() {
List<StringBuilder> conditions = new ArrayList<>();
String REGISTRY_DESIGNATION_FORMAT = "registryDesignation = '%s'";
String REGISTRY_INSTRUMENT_TYPE_FORMAT = "registryInstrumentType = '%s'";
String REGISTRY_CAPACITY_FORMAT = "registryCapacity = '%s'";
String REGISTRY_UNIT_FORMAT = "registryUnit = '%s'";
if (prefix != null) {
REGISTRY_DESIGNATION_FORMAT = "%s.%s".formatted(prefix, REGISTRY_DESIGNATION_FORMAT);
REGISTRY_INSTRUMENT_TYPE_FORMAT = "%s.%s".formatted(prefix, REGISTRY_INSTRUMENT_TYPE_FORMAT);
REGISTRY_CAPACITY_FORMAT = "%s.%s".formatted(prefix, REGISTRY_CAPACITY_FORMAT);
REGISTRY_UNIT_FORMAT = "%s.%s".formatted(prefix, REGISTRY_UNIT_FORMAT);
}
for (RegistryTradingParams tradingParams : this.registryTradingParams) {
StringBuilder condition = new StringBuilder();
boolean wasAddedCondition = false;
if (tradingParams.registryDesignation() != null) {
condition.append(String.format(REGISTRY_DESIGNATION_FORMAT, tradingParams.registryDesignation().getKey()));
condition.append(String.format("registryDesignation = '%s'", tradingParams.registryDesignation().getKey()));
wasAddedCondition = true;
}
if (wasAddedCondition) {
@ -46,7 +35,7 @@ public class RegistryCodeSqlBuilder {
wasAddedCondition = false;
}
if (tradingParams.registryInstrumentType() != null) {
condition.append(String.format(REGISTRY_INSTRUMENT_TYPE_FORMAT, tradingParams.registryInstrumentType().getKey()));
condition.append(String.format("registryInstrumentType = '%s'", tradingParams.registryInstrumentType().getKey()));
wasAddedCondition = true;
}
if (wasAddedCondition) {
@ -54,7 +43,7 @@ public class RegistryCodeSqlBuilder {
wasAddedCondition = false;
}
if (tradingParams.registryCapacity() != null) {
condition.append(String.format(REGISTRY_CAPACITY_FORMAT, tradingParams.registryCapacity().getKey()));
condition.append(String.format("registryCapacity = '%s'", tradingParams.registryCapacity().getKey()));
wasAddedCondition = true;
}
if (wasAddedCondition) {
@ -62,7 +51,7 @@ public class RegistryCodeSqlBuilder {
wasAddedCondition = false;
}
if (tradingParams.registryUnit() != null) {
condition.append(String.format(REGISTRY_UNIT_FORMAT, tradingParams.registryUnit().getKey()));
condition.append(String.format("registryUnit = '%s'", tradingParams.registryUnit().getKey()));
wasAddedCondition = true;
}
if (!wasAddedCondition) {
@ -73,31 +62,22 @@ public class RegistryCodeSqlBuilder {
return conditions.stream().collect(Collectors.joining(") or (", "(", ")"));
}
public ImdgPredicate buildPredicate(ImdgPredicateBuilder pb) {
List<ImdgPredicate> conditionsOr = new ArrayList<>();
String REGISTRY_DESIGNATION_FIELD = "registryDesignation";
String REGISTRY_INSTRUMENT_TYPE_FIELD = "registryInstrumentType";
String REGISTRY_CAPACITY_FIELD = "registryCapacity";
String REGISTRY_UNIT_FIELD = "registryUnit";
if (prefix != null) {
REGISTRY_DESIGNATION_FIELD = "%s.%s".formatted(prefix, REGISTRY_DESIGNATION_FIELD);
REGISTRY_INSTRUMENT_TYPE_FIELD = "%s.%s".formatted(prefix, REGISTRY_INSTRUMENT_TYPE_FIELD);
REGISTRY_CAPACITY_FIELD = "%s.%s".formatted(prefix, REGISTRY_CAPACITY_FIELD);
REGISTRY_UNIT_FIELD = "%s.%s".formatted(prefix, REGISTRY_UNIT_FIELD);
}
for (RegistryTradingParams tradingParams : this.registryTradingParams) {
List<ImdgPredicate> conditionsAnd = new ArrayList<>(4);
if (tradingParams.registryDesignation() != null) {
conditionsAnd.add(pb.equals(REGISTRY_DESIGNATION_FIELD, tradingParams.registryDesignation().getKey()));
conditionsAnd.add(pb.equals("registryDesignation", tradingParams.registryDesignation().getKey()));
}
if (tradingParams.registryInstrumentType() != null) {
conditionsAnd.add(pb.equals(REGISTRY_INSTRUMENT_TYPE_FIELD, tradingParams.registryInstrumentType().getKey()));
conditionsAnd.add(pb.equals("registryInstrumentType", tradingParams.registryInstrumentType().getKey()));
}
if (tradingParams.registryCapacity() != null) {
conditionsAnd.add(pb.equals(REGISTRY_CAPACITY_FIELD, tradingParams.registryCapacity().getKey()));
conditionsAnd.add(pb.equals("registryCapacity", tradingParams.registryCapacity().getKey()));
}
if (tradingParams.registryUnit() != null) {
conditionsAnd.add(pb.equals(REGISTRY_UNIT_FIELD, tradingParams.registryUnit().getKey()));
conditionsAnd.add(pb.equals("registryUnit", tradingParams.registryUnit().getKey()));
}
if (conditionsAnd.size() == 1)
conditionsOr.add(conditionsAnd.get(0));
@ -108,8 +88,4 @@ public class RegistryCodeSqlBuilder {
if (conditionsOr.size() == 1) return conditionsOr.get(0);
return pb.or(conditionsOr.toArray(new ImdgPredicate[conditionsOr.size()]));
}
public RegistryCodeSqlBuilder prefix(String prefix) {
this.prefix = prefix;
return this;
}
}

View file

@ -22,9 +22,6 @@ public class KafkaConsumerFactory {
props.put("enable.auto.commit", kafkaSettings.getEnableAutoCommit().toString());
props.put("session.timeout.ms", kafkaSettings.getSessionTimeoutMs().toString());
props.put("auto.offset.reset", kafkaSettings.getAutoOffsetReset());
if (kafkaSettings.getMaxPollIntervalMs() != null) {
props.put("max.poll.interval.ms", String.valueOf(kafkaSettings.getMaxPollIntervalMs()));
}
props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
return new KafkaConsumer<>(props);
@ -39,9 +36,6 @@ public class KafkaConsumerFactory {
}
props.put("enable.auto.commit", kafkaSettings.getEnableAutoCommit().toString());
props.put("session.timeout.ms", kafkaSettings.getSessionTimeoutMs().toString());
if (kafkaSettings.getMaxPollIntervalMs() != null) {
props.put("max.poll.interval.ms", String.valueOf(kafkaSettings.getMaxPollIntervalMs()));
}
props.put("auto.offset.reset", kafkaSettings.getAutoOffsetReset());
props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
// props.put("value.deserializer", JsonDeserializer.class.getName());

View file

@ -6,7 +6,6 @@ public class KafkaConsumerSettings {
private Boolean enableAutoCommit;
private Integer sessionTimeoutMs;
private String autoOffsetReset;
private Integer maxPollIntervalMs = null;
public String getBootstrapServers() {
@ -48,13 +47,4 @@ public class KafkaConsumerSettings {
public void setAutoOffsetReset(String autoOffsetReset) {
this.autoOffsetReset = autoOffsetReset;
}
public Integer getMaxPollIntervalMs() {
return maxPollIntervalMs;
}
public void setManuallyMaxPollIntervalMs(Integer maxPollIntervalMs) {
this.maxPollIntervalMs = maxPollIntervalMs;
}
}

View file

@ -0,0 +1,8 @@
package ru.spcex.clearing.platform.messaging.service;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
public interface IKafkaRequestStatusProcessor {
default void updateRequestStatus(BaseRequest<?> o, Object response) {}
default void updateRequestStatusOnError(BaseRequest<?> o, Object response) {}
}