Merge remote-tracking branch 'origin/dev' into psemenkov

# Conflicts:
#	clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/config/KafkaTestConfig.java
#	clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/AbstractServiceTest.java
This commit is contained in:
psemenkov 2023-02-03 17:12:01 +03:00
commit 225adec1f0
67 changed files with 1982 additions and 261 deletions

View file

@ -24,12 +24,13 @@ import ru.spcex.clearing.backendapi.security.KeycloakUtils;
import ru.spcex.clearing.backendapi.service.IOperator;
import ru.spcex.clearing.backendapi.service.IStateLoader;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.platform.enumeration.Task;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.enumeration.IMessageResolver;
import java.time.Instant;
import java.util.Collection;
import java.util.Map;
import java.util.concurrent.ExecutionException;
@ -81,19 +82,36 @@ public class LauncherController extends AbstractQueueController {
}
launcherCommand.setTask(dictionaryName);
launcherCommand.setUserId(user.getId());
// пока здесь, это требуется для сохранения истории
saveLauncher(dictionaryName, user.getId());
//топики ограничиваются наличием в taskDictionary
//подписываются на разные топики в разных модулях, см. ru.spcex.platform.enumeration.Task#topic
return processRequest("launcher-" + dictionaryName, launcherCommand);
return processRequest(Consts.LAUNCHER_NEW, launcherCommand);
}
private void saveLauncher(String taskCode, Long userId) {
Launcher launcherCommandHistory = new Launcher();
launcherCommandHistory.setTask(taskCode);
launcherCommandHistory.setSenderId(userId);
launcherCommandHistory.setCreated(Instant.now());
launcherImdg.insert(launcherCommandHistory);
@ApiOperation(value = "create specific launcher.")
@ApiResponses(value = {@ApiResponse(code = 200, message = "OK", response = CudResponse.class), @ApiResponse(code = 400, message = "Ошибка валидации", response = BasicSpcexResponse.class)})
@RequestMapping(value = "/specific", method = RequestMethod.POST, produces = MediaType.APPLICATION_JSON_VALUE)
@ResponseBody
public CudResponse addSpecific(@ApiParam(value = "Параметры команды в JSON формате.", required = true)
@RequestBody LauncherNew launcherNew) throws ExecutionException, InterruptedException {
AbstractDictionary taskEnum = taskDictionary.getSingleObjectByFieldValues(Map.of("code", launcherNew.getTask()));
if (taskEnum == null) {
throw new NotFound404Exception("task dictionary element with code '" + launcherNew.getTask() + "'");
}
if (!Task.startOfClearing.getKey().equals(taskEnum.getCode())) {
throw new IllegalStateException(String.format("Task %s not support request with body", taskEnum.getCode()));
}
LauncherNew launcherCommand = new LauncherNew();
Authentication authentication = SecurityContextHolder.getContext().getAuthentication();
String username = KeycloakUtils.getUserNameFromAuthentication(authentication);
User user = userImdg.getSingleObjectByFieldValues(Map.of("identifier", username));
if (user == null) {
throw new IllegalStateException("cannot obtain userId from logged in user " + username);
}
launcherCommand.setTask(launcherNew.getTask());
launcherCommand.setUserId(user.getId());
launcherCommand.setCompanyId(launcherNew.getCompanyId());
launcherCommand.setSecurityId(launcherNew.getSecurityId());
return processRequest(Consts.LAUNCHER_NEW, launcherCommand);
}
@Autowired

View file

@ -6,9 +6,7 @@ import io.swagger.annotations.ApiModelProperty;
import ru.spcex.clearing.backendapi.domain.actions.IAction;
import ru.spcex.clearing.backendapi.errors.BackEndError;
import ru.spcex.clearing.platform.messaging.domain.ActionType;
import ru.spcex.clearing.platform.messaging.domain.cud.reports.ReportWithPeriodRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
import ru.spcex.platform.enumeration.Task;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import java.util.Collection;
@ -19,24 +17,26 @@ public class LauncherNew implements IAction<Object> {
@ApiModelProperty(value = "Идентификатор единоличного исполнительного органа", example = "ABCD")
@JsonProperty
private String task;
@ApiModelProperty(value = "Идентификатор инициатора", example = "1000")
@JsonProperty
private Long companyId;
@ApiModelProperty(value = "Идентификатор инструмента", example = "1000")
@JsonProperty
private Long securityId;
@JsonIgnore
private Long userId;
@Override
public Object toRequest() {
if (Task.createReport_GREP.equalsByKey(task)) {
ReportWithPeriodRequest reportCommand = new ReportWithPeriodRequest();
reportCommand.setReportId(task);
return reportCommand;
} else {
LauncherCommandRequest taskRunnerCommandRequest = new LauncherCommandRequest();
taskRunnerCommandRequest.setTaskName(task);
taskRunnerCommandRequest.setCompanyId(companyId);
taskRunnerCommandRequest.setSecurityId(securityId);
taskRunnerCommandRequest.setUserId(userId);
return taskRunnerCommandRequest;
}
}
@ApiModelProperty(hidden = true)
@Override
public ActionType getActionType() {
return ActionType.NEW;
@ -65,4 +65,20 @@ public class LauncherNew implements IAction<Object> {
public void setUserId(Long userId) {
this.userId = userId;
}
public Long getCompanyId() {
return companyId;
}
public void setCompanyId(Long companyId) {
this.companyId = companyId;
}
public Long getSecurityId() {
return securityId;
}
public void setSecurityId(Long securityId) {
this.securityId = securityId;
}
}

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="2.4.0.7">
<meta version="2.4.0.9">
<!-- _xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" _xsi:noNamespaceSchemaLocation="file:///E:/d/projects/meta/from/meta.xsd" -->
<!--Здесь словари-->
<enums>
@ -713,8 +713,8 @@
</actions>
</bankAccount>
<informationAccount name="Информационные счета" destination="securities/information-accounts" class="ru.clearing.classes.statics.data.account.InformationAccount" logUpdates="true" table="information_account">
<accountId type="1" name="Счет" shortname="Счет" searchable="true" sortable="true" visible="true" link="account" linkCode="account"/>
<clearingAccountId type="1" name="Счета" shortname="Счет" searchable="true" sortable="true" visible="true" link="account" linkCode="account"/>
<accountId type="1" name="Информационный счет" shortname="Информационный счет" searchable="true" sortable="true" visible="true" link="account" linkCode="account"/>
<clearingAccountId type="1" name="Аналитический счет" shortname="Аналитический счет" searchable="true" sortable="true" visible="true" link="account" linkCode="account"/>
<companyId type="1" name="Компания" shortname="Компания" searchable="true" sortable="true" visible="true" link="company" linkCode="shortName"/>
<id type="1" name="Идентификатор записи" shortname="ID" searchable="true" sortable="true"/>
</informationAccount>
@ -752,7 +752,7 @@
<instrumentType type="12" name="Наименование типа инструмента" shortname="Тип инструмента" searchable="true" sortable="true" visible="true" link="instrumentType" extends="security"/>
<fullName type="2" name="Полное наименование инструмента" shortname="Наименование" searchable="true" sortable="true" length="255" visible="true" extends="security"/>
<securitySymbol type="2" name="Код инструмента" shortname="Код" searchable="true" sortable="true" length="255" visible="true" extends="security"/>
<termType type="12" name="Наименование вида инструмента" shortname="Вид инструмента" searchable="true" sortable="true" visible="true" link="termType"/>
<termType type="12" name="Наименование вида инструмента" shortname="Вид инструмента" searchable="true" sortable="true" visible="false" link="termType" ignore="true"/>
<lotSize field="securityId" type="11" name="Размер лота" shortname="Размер лота" searchable="true" sortable="true" visible="true" linkKeyCode="securityId" linkCode="lotSize" link="listing"/>
<id type="1" name="Идентификатор записи" shortname="ID" searchable="true" sortable="true"/>
<actions>
@ -900,20 +900,20 @@
<tradingDate type="6" name="Дата заключения сделки" shortname="Дата заключения сделки" visible="true" searchable="true" sortable="true"/>
<accountId type="1" name="Торговый счет" shortname="Счет" visible="true" searchable="true" sortable="true" link="account" linkCode="account"/>
<market type="12" name="Секция финансового инструмента" shortname="Секция" visible="true" searchable="true" sortable="true" link="market"/>
<price field="financialProduct.interestRate.value" type="10" name="Ставка по депозиту" shortname="Ставка, %" visible="true" searchable="true" sortable="true"/>
<price type="10" name="Ставка по депозиту" shortname="Ставка, %" visible="true" searchable="true" sortable="true"/>
<lots type="11" name="Количество лотов" shortname="Лоты" visible="true" searchable="true" sortable="true"/>
<quantity type="11" name="Количество штук" shortname="Штуки" visible="false" searchable="true" sortable="true"/>
<firstLegAmount type="11" name="Объем сделки" shortname="Объем" visible="true" searchable="true" sortable="true"/>
<secondLegAmount type="11" name="Объем возврата" shortname="Объем возврата" visible="false" searchable="true" sortable="true"/>
<interestAmount type="11" name="Объем процентов" shortname="Проценты" visible="false" searchable="true" sortable="true"/>
<side type="12" name="Направление сделки" shortname="Направление" visible="true" searchable="true" sortable="true" link="moneyFlowSide"/>
<settlementCurrency type="12" name="Валюта расчетов по инструменту" shortname="Валюта" visible="true" searchable="true" sortable="true" link="currencyCode" linkCode="currencyCode"/>
<companyId field="company.id" type="1" name="Название компании" shortname="Компания" visible="true" searchable="true" sortable="true" link="company" linkCode="shortName"/>
<settlementCurrency type="12" name="Валюта расчетов по инструменту" shortname="Валюта" visible="true" searchable="true" sortable="true" link="currencyCode" linkCode="code"/>
<companyId type="1" name="Название компании" shortname="Компания" visible="true" searchable="true" sortable="true" link="company" linkCode="shortName"/>
<duration type="1" name="Срок, дней" shortname="Срок" visible="true" searchable="true" sortable="true"/>
<firstLegSettlementDate field="firstLeg.settlementDate" type="6" name="Дата размещения" shortname="Дата размещения" visible="true" searchable="true" sortable="true"/>
<secondLegSettlementDate field="secondLeg.settlementDate" type="6" name="Дата возврата" shortname="Дата возврата" visible="true" searchable="true" sortable="true"/>
<firstLegSettlementCode field="firstLeg.settlementCode" type="6" name="Код расчетов при размещении" shortname="Код расчетов при размещении" visible="false" searchable="true" sortable="true"/>
<secondLegSettlementCode field="secondLeg.settlementCode" type="6" name="Код расчетов при возврате" shortname="Код расчетов" visible="true" searchable="true" sortable="true"/>
<firstLegSettlementDate type="6" name="Дата размещения" shortname="Дата размещения" visible="true" searchable="true" sortable="true"/>
<secondLegSettlementDate type="6" name="Дата возврата" shortname="Дата возврата" visible="true" searchable="true" sortable="true"/>
<firstLegSettlementCode type="6" name="Код расчетов при размещении" shortname="Код расчетов при размещении" visible="false" searchable="true" sortable="true" ignore="true"/>
<secondLegSettlementCode type="6" name="Код расчетов при возврате" shortname="Код расчетов" visible="false" searchable="true" sortable="true" ignore="true"/>
<securityFullName type="2" length="255" name="Наименование инструмента" shortname="Инструмент" visible="true" searchable="true" sortable="true"/>
<securitySymbol type="2" length="255" name="Код инструмента в Торговой Системе" shortname="Код инструмента" visible="true" searchable="true" sortable="true"/>
<securityId type="1" name="Финансовый инструмент" shortname="Код биржевого инструмента / товара" searchable="true" sortable="true" link="moneyMarketSecurity" linkCode="securitySymbol" ignore="true"/>
@ -921,8 +921,8 @@
<coverageStatus type="12" name="Cтатус достаточности обеспечения" shortname="Обеспеченность" searchable="true" sortable="true" link="allowed"/>
<sessionId type="1" name="Наименование сессии" shortname="Сессия" visible="true" searchable="true" sortable="true" link="moneyMarketSession"/>
<id type="1" name="Идентификатор записи" shortname="ID" visible="false" searchable="true" sortable="true"/>
<createdAt type="5" name="Время регистрации сделки" shortname="Время сделки" visible="false" searchable="true" sortable="true"/>
<updatedAt type="5" name="Время изменения сделки" shortname="Время изменения" visible="false" searchable="true" sortable="true"/>
<createdAt field="created" type="5" name="Время регистрации сделки" shortname="Время сделки" visible="false" searchable="true" sortable="true"/>
<updatedAt field="updated" type="5" name="Время изменения сделки" shortname="Время изменения" visible="false" searchable="true" sortable="true"/>
<clearingDate type="6" name="Дата клиринга" shortname="Дата клиринга" visible="false" searchable="true" sortable="true"/>
</executionDeposit>
<dealRegister name="Реестр сделок" destination="deal-registers" class="ru.clearing.classes.statics.data.register.DealRegister" table="deal_register">

View file

@ -3,13 +3,20 @@ package ru.spcex.clearing.balance.config;
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.producer.Producer;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.beans.factory.config.ConfigurableBeanFactory;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Scope;
import ru.spcex.clearing.balance.config.element.BalanceServiceSettings;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.config.KafkaConsumerFactory;
import ru.spcex.clearing.platform.messaging.config.KafkaProducerFactory;
import ru.spcex.clearing.platform.messaging.service.RequestInfo;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgId;
import ru.spcex.platform.imdg.api.ImdgProvider;
@Configuration
public class KafkaConfig {
@ -21,8 +28,24 @@ public class KafkaConfig {
}
@Autowired
@Bean
@Bean("kafkaProducer")
public Producer<String, Object> createProducer(BalanceServiceSettings settings) {
return KafkaProducerFactory.producer(settings.getKafkaProducer());
}
@Autowired
@Bean
public KafkaSender kafkaSender(@Qualifier("kafkaProducer") Producer<String, Object> kafkaProducer,
ImdgProvider imdgProvider) {
ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator();
return KafkaSender
.setup()
.producer(kafkaProducer)
.idGenerator(imdgIdGenerator::nextId)
.imdgProvider(s -> {
Imdg<RequestInfo> imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class);
return imdg::insert;
})
.build();
}
}

View file

@ -1,31 +0,0 @@
package ru.spcex.clearing.balance.config;
import org.apache.kafka.clients.producer.Producer;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.service.RequestInfo;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgId;
import ru.spcex.platform.imdg.api.ImdgProvider;
@Configuration
public class KafkaSenderConfig {
@Autowired
@Bean
public KafkaSender kafkaSender(Producer<String, Object> kafkaProducer, ImdgProvider imdgProvider) {
ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator();
return KafkaSender
.setup()
.producer(kafkaProducer)
.idGenerator(imdgIdGenerator::nextId)
.imdgProvider(s -> {
Imdg<RequestInfo> imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class);
return imdg::insert;
})
.build();
}
}

View file

@ -1,17 +1,25 @@
package ru.spcex.clearing.balance.service;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.stereotype.Component;
import ru.clearing.classes.statics.data.account.Account;
import ru.clearing.classes.statics.data.account.AccountBalance;
import ru.clearing.classes.statics.data.company.Company;
import ru.clearing.classes.statics.data.generated.ClearingMemberCategory;
import ru.spcex.clearing.balance.validation.AccountBalanceValidation;
import ru.spcex.clearing.balance.validation.ValidationStored;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.cud.balance.AccountBalanceClearingRequest;
import ru.spcex.platform.enumeration.AccountType;
import ru.spcex.platform.enumeration.BalanceAccountType;
import ru.spcex.platform.enumeration.ClearingCategory;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.enumeration.IEnumKey;
import ru.spcex.platform.utils.number.BigDecimalUtil;
import ru.spcex.platform.utils.validation.IValidator;
import java.math.BigDecimal;
@ -23,15 +31,18 @@ import java.util.function.Function;
@Component
public class AccountBalanceService {
private final Logger log = LoggerFactory.getLogger(getClass());
private final ImdgProvider imdgProvider;
private final Imdg<AccountBalance> accountBalanceImdg;
private final Function<AccountBalanceValidation, IValidator> validationFactory;
private final Imdg<ClearingMemberCategory> clearingCategoryImdg;
public AccountBalanceService(ImdgProvider imdgProvider,
@Qualifier("accountBalanceValidator") Function<AccountBalanceValidation, IValidator> validationFactory) {
this.imdgProvider = imdgProvider;
this.validationFactory = validationFactory;
this.accountBalanceImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_AccountBalance, AccountBalance.class);
this.clearingCategoryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ClearingMemberCategory, ClearingMemberCategory.class);
}
public AccountResult createAccountBalance(Long addresseeId, Long accountId, BigDecimal amount,
@ -85,4 +96,101 @@ public class AccountBalanceService {
if (b == null) return a;
return a.add(b);
}
public void updateAccountBalanceByClearing(AccountBalanceClearingRequest req) {
ClearingCategory category = getClearingCategoryByCompanyId(req.getCompanyId());
req.setFirstLegAmount(BigDecimalUtil.safeBD(req.getFirstLegAmount()));
if (category.equals(ClearingCategory.I)) {
updateAccountCategoryIClrn(req.getAccountId(), req.getCompanyId(), req.getFirstLegAmount());
updateAccountCategoryITran(req.getFirstLegAmount());
} else if (category.equals(ClearingCategory.V)) {
updateAccountCategoryVInfo(req.getAccountId(), req.getCompanyId(), req.getFirstLegAmount());
updateAccountCategoryVAnlt(req.getFirstLegAmount());
} else if (category.equals(ClearingCategory.B)) {
updateAccountCategoryBClrn(req.getAccountId(), req.getCompanyId(), req.getFirstLegAmount());
}
}
private ClearingCategory getClearingCategoryByCompanyId(Long companyId) {
ClearingMemberCategory category = clearingCategoryImdg.getSingleObjectByFieldValues(
Map.of("companyId", companyId));
return IEnumKey.getEnumByKey(ClearingCategory.class,
category.getClearingMemberCategory());
}
private void updateAccountCategoryIClrn(Long accountId, Long companyId, BigDecimal firstLegAmount) {
AccountBalance accountBalance = loadAccountBalance(accountId, companyId, AccountType.Clrn);
log.trace("updating AccountBalance {} I CLRN", accountBalance.getId());
BigDecimal previousFreeBalanceAmount = BigDecimalUtil.safeBD(accountBalance.getFreeBalanceAmount());
BigDecimal previousChangeBalanceAmount = BigDecimalUtil.safeBD(accountBalance.getChangeBalanceAmount());
BigDecimal previousBalanceAmount = BigDecimalUtil.safeBD(accountBalance.getBalanceAmount());
accountBalance.setFreeBalanceAmount(previousFreeBalanceAmount.subtract(firstLegAmount));
accountBalance.setChangeBalanceAmount(previousChangeBalanceAmount.subtract(firstLegAmount));
accountBalance.setBalanceAmount(previousBalanceAmount.subtract(firstLegAmount));
accountBalance.setUpdated(Instant.now());
log.trace("updateAccountCategoryIClrn: {}", accountBalance.getId());
accountBalanceImdg.update(accountBalance);
}
private void updateAccountCategoryITran(BigDecimal firstLegAmount) {
AccountBalance accountBalance = loadAccountBalance(AccountType.Tran);
log.trace("updating AccountBalance {} I TRAN", accountBalance.getId());
BigDecimal previousDebitAmount = BigDecimalUtil.safeBD(accountBalance.getDebitAmount());
accountBalance.setDebitAmount(previousDebitAmount.add(firstLegAmount.abs()));
accountBalance.setUpdated(Instant.now());
log.trace("updateAccountCategoryITran: {}", accountBalance.getId());
accountBalanceImdg.update(accountBalance);
}
private void updateAccountCategoryVInfo(Long accountId, Long companyId, BigDecimal firstLegAmount) {
AccountBalance accountBalance = loadAccountBalance(accountId, companyId, AccountType.Info);
log.trace("updating AccountBalance {} V INFO", accountBalance.getId());
BigDecimal previousFreeBalanceAmount = BigDecimalUtil.safeBD(accountBalance.getFreeBalanceAmount());
BigDecimal changeBalanceAmount = BigDecimalUtil.safeBD(accountBalance.getChangeBalanceAmount());
BigDecimal balanceAmount = BigDecimalUtil.safeBD(accountBalance.getBalanceAmount());
accountBalance.setFreeBalanceAmount(previousFreeBalanceAmount.subtract(firstLegAmount));
accountBalance.setChangeBalanceAmount(changeBalanceAmount.subtract(firstLegAmount));
accountBalance.setBalanceAmount(balanceAmount.subtract(firstLegAmount));
accountBalance.setUpdated(Instant.now());
log.trace("updateAccountCategoryVInfo: {}", accountBalance.getId());
accountBalanceImdg.update(accountBalance);
}
private void updateAccountCategoryVAnlt(BigDecimal firstLegAmount) {
AccountBalance accountBalance = loadAccountBalance(AccountType.Anlt);
log.trace("updating AccountBalance {} V ANLT", accountBalance.getId());
BigDecimal previousFreeBalanceAmount = BigDecimalUtil.safeBD(accountBalance.getFreeBalanceAmount());
BigDecimal previousChangeBalanceAmount = BigDecimalUtil.safeBD(accountBalance.getChangeBalanceAmount());
BigDecimal previousBalanceAmount = BigDecimalUtil.safeBD(accountBalance.getBalanceAmount());
accountBalance.setFreeBalanceAmount(previousFreeBalanceAmount.subtract(firstLegAmount));
accountBalance.setChangeBalanceAmount(previousChangeBalanceAmount.subtract(firstLegAmount));
accountBalance.setBalanceAmount(previousBalanceAmount.subtract(firstLegAmount));
accountBalance.setUpdated(Instant.now());
log.trace("updateAccountCategoryVAnlt: {}", accountBalance.getId());
accountBalanceImdg.update(accountBalance);
}
private void updateAccountCategoryBClrn(Long accountId, Long companyId, BigDecimal firstLegAmount) {
AccountBalance accountBalance = loadAccountBalance(accountId, companyId, AccountType.Clrn);
log.trace("updating AccountBalance {} B CLRN", accountBalance.getId());
BigDecimal debitAmount = BigDecimalUtil.safeBD(accountBalance.getDebitAmount());
accountBalance.setDebitAmount(plus(debitAmount, firstLegAmount.abs()));
accountBalance.setUpdated(Instant.now());
log.trace("updateAccountCategoryBClrn: {}", accountBalance.getId());
accountBalanceImdg.update(accountBalance);
}
private AccountBalance loadAccountBalance(Long accountId, Long companyId, AccountType type) {
return accountBalanceImdg.getSingleObjectByFieldValues(
Map.of("accountId", accountId, "companyId", companyId,
"accountType", type.getKey())
);
}
private AccountBalance loadAccountBalance(AccountType type) {
return accountBalanceImdg.getSingleObjectByFieldValues(
Map.of("companyId", 1L,
"accountType", type.getKey())
);
}
}

View file

@ -15,8 +15,10 @@ import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.account.sdf01.AccountSdf01Request;
import ru.spcex.clearing.platform.messaging.domain.cud.account.sdf01.AccountSdfRequestPart;
import ru.spcex.clearing.platform.messaging.domain.cud.balance.AccountBalanceClearingRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.balance.ExportToFileRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonIdRequest;
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.platform.classes.base.interfaces.WithAccount;
@ -38,14 +40,16 @@ public class StatementService extends QueueConsumer implements InitializingBean
private final KafkaSender kafkaReqProducer;
private final Map<SdfTable, Imdg<? extends WithAccount>> sdfImdgs;
private final Map<SdfTable, AbstractExecutor<? extends WithAccount>> executorsMap;
private final AccountBalanceService accountBalanceService;
@Autowired
public StatementService(Consumer<String, Object> kafkaQueue,
ImdgProvider imdgProvider,
KafkaSender kafkaReqProducer,
@Qualifier("sdfExecutors") Map<SdfTable, AbstractExecutor<? extends WithAccount>> executorsMap) {
@Qualifier("sdfExecutors") Map<SdfTable, AbstractExecutor<? extends WithAccount>> executorsMap, AccountBalanceService accountBalanceService) {
super(kafkaQueue);
this.imdgProvider = imdgProvider;
this.accountBalanceService = accountBalanceService;
this.sdfImdgs = new EnumMap<>(SdfTable.class);
this.kafkaReqProducer = kafkaReqProducer;
this.executorsMap = executorsMap;
@ -59,9 +63,22 @@ public class StatementService extends QueueConsumer implements InitializingBean
callback(StatementRequest.class)
.setConsumer(this::process)
.forDestination(Consts.STATEMENT_PROCESS, callbacks::put);
callback(AccountBalanceClearingRequest.class)
.setConsumer(this::accountBalanceClearingUpdate)
.forDestination(Consts.BALANCE_ACCOUNT_UPDATE, callbacks::put);
init();
}
private void accountBalanceClearingUpdate(BaseRequest<AccountBalanceClearingRequest> updateAccBalanceReq) {
AccountBalanceClearingRequest payload = updateAccBalanceReq.getRequestPayload();
log.debug("updating account balance for clearing baseRequest.id = {}; accountId = {}, companyId = {}, firstLegAmount = {}", updateAccBalanceReq.getId(),
payload.getAccountId(), payload.getCompanyId(), payload.getFirstLegAmount());
accountBalanceService.updateAccountBalanceByClearing(payload);
CommonIdRequest commonIdRequest = new CommonIdRequest();
commonIdRequest.setId(updateAccBalanceReq.getId());
kafkaReqProducer.sendRequestToQueue(Consts.CONTINUE_CLEARING, commonIdRequest);
}
private void process(BaseRequest<StatementRequest> systemRequest) {
StatementRequest statementRequest = systemRequest.getRequestPayload();
Collection<? extends WithAccount> sdfGroup;

View file

@ -51,6 +51,10 @@
<artifactId>assertj-core</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.mockito</groupId>
<artifactId>mockito-core</artifactId>
</dependency>
</dependencies>
<build>
<resources>

View file

@ -0,0 +1,34 @@
package ru.spcex.clearing.config;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import ru.clearing.classes.statics.data.generated.ClearingMemberCategory;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.platform.enumeration.ClearingCategory;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.utils.enumeration.IEnumKey;
import java.util.Map;
import java.util.function.Function;
@Configuration
public class SortingConfig {
private final Imdg<ClearingMemberCategory> clearingCategoryImdg;
@Autowired
public SortingConfig(ImdgProvider imdgProvider) {
this.clearingCategoryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ClearingMemberCategory, ClearingMemberCategory.class);
}
@Bean
public Function<Long, ClearingCategory> clearingMemberCategoryProvider() {
return companyId -> {
ClearingMemberCategory category = clearingCategoryImdg.getSingleObjectByFieldValues(
Map.of("companyId", companyId));
return IEnumKey.getEnumByKey(ClearingCategory.class,
category.getClearingMemberCategory());
};
}
}

View file

@ -0,0 +1,65 @@
package ru.spcex.clearing.config;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import ru.clearing.classes.statics.data.company.Company;
import ru.clearing.classes.statics.data.company.relation.Relation;
import ru.clearing.classes.statics.data.execution.ExecutionDeposit;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.service.validation.ClrngValidationStored;
import ru.spcex.clearing.service.validation.ExecutionDepositValidationRule;
import ru.spcex.platform.classes.base.SpcexObjectBase;
import ru.spcex.platform.enumeration.ClearingCategory;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.imdg.validation.ImdgValidationContext;
import ru.spcex.platform.utils.validation.IValidator;
import ru.spcex.platform.utils.validation.ValidatorImpl;
import java.util.HashMap;
import java.util.Map;
import java.util.function.BiConsumer;
import java.util.function.BiFunction;
@Configuration
public class ValidationConfig {
private final Map<String, Imdg<? extends SpcexObjectBase>> imdgs;
public ValidationConfig(ImdgProvider imdgProvider) {
this.imdgs = new HashMap<>();
BiConsumer<String, Class<? extends SpcexObjectBase>> addImdg = (s, aClass) -> imdgs.put(s, imdgProvider.getImdg(s, aClass));
addImdg.accept(IMDGDistributedNames.Map_Relation, Relation.class);
addImdg.accept(IMDGDistributedNames.Map_Company, Company.class);
addImdg.accept(IMDGDistributedNames.Map_Account, Company.class);
addImdg.accept(IMDGDistributedNames.Map_AccountBalance, Company.class);
}
/**
* чтобы во всех валидаторах был один экземпляр Imdg
*/
private Imdg<?> getImdg(String key) {
return imdgs.get(key);
}
private final BiConsumer<ImdgValidationContext<?>, String> addImdg = (context, s)
-> context.addImdg(s, getImdg(s));
@Bean("executionDepositValidator")
public BiFunction<ClearingCategory, ExecutionDeposit, IValidator> executionDepositValidator() {
return (category, execDeposit) -> {
ImdgValidationContext<ExecutionDeposit> context = new ImdgValidationContext<>();
context.setValidatedObject(execDeposit);
context.storeObject(ClrngValidationStored.categoryPredefined, category);
addImdg.accept(context, IMDGDistributedNames.Map_Relation);
addImdg.accept(context, IMDGDistributedNames.Map_Company);
addImdg.accept(context, IMDGDistributedNames.Map_Account);
addImdg.accept(context, IMDGDistributedNames.Map_AccountBalance);
return new ValidatorImpl<>(context,
ExecutionDepositValidationRule.ClearingIsAllowed,
ExecutionDepositValidationRule.AccountIsNotBlocked,
ExecutionDepositValidationRule.CompanyIsNotBlocked,
ExecutionDepositValidationRule.FinancialObligationSecurity);
};
}
}

View file

@ -0,0 +1,21 @@
package ru.spcex.clearing.error;
import ru.spcex.platform.utils.enumeration.IEnumId;
public enum ClearingErrorInternal implements IEnumId {
ClearingNotAllowed(1L),
AccountNotActive(2L),
CompanyNotActive(3L),
FinancialObligationNotSatisfied(4L),
;
private final Long id;
ClearingErrorInternal(Long id) {
this.id = id;
}
@Override
public Long getId() {
return id;
}
}

View file

@ -0,0 +1,234 @@
package ru.spcex.clearing.service;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.execution.ExecutionDeposit;
import ru.clearing.classes.statics.data.generated.ClearingMemberCategory;
import ru.clearing.classes.statics.data.liabilities.LiabilitiesClaimsAssets;
import ru.clearing.classes.statics.data.liabilities.LiabilitiesClaimsMoney;
import ru.clearing.classes.statics.data.payment.PaymentInstruction;
import ru.spcex.clearing.error.ClearingErrorInternal;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.balance.AccountBalanceClearingRequest;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.clearing.service.builder.LiabilitiesClaimsAssetsCreator;
import ru.spcex.clearing.service.builder.LiabilitiesClaimsMoneyCreator;
import ru.spcex.clearing.service.builder.PaymentInstructionCreator;
import ru.spcex.clearing.service.order.ExecutionDepositSorter;
import ru.spcex.platform.enumeration.Allowed;
import ru.spcex.platform.enumeration.ClearingCategory;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgId;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.log.ExceptionUtils;
import ru.spcex.platform.utils.validation.IValidator;
import java.time.Instant;
import java.util.*;
import java.util.function.BiFunction;
@Service
public class Clearing {
private final Logger log = LoggerFactory.getLogger(getClass());
private final Imdg<ExecutionDeposit> executionDepositImdg;
private final Imdg<LiabilitiesClaimsAssets> liabilitiesClaimsAssetsImdg;
private final Imdg<LiabilitiesClaimsMoney> liabilitiesClaimsMoneyImdg;
private final Imdg<ClearingMemberCategory> clearingCategoryImdg;
private final ExecutionDepositSorter sorter;
private final BiFunction<ClearingCategory, ExecutionDeposit, IValidator> validation;
private final PaymentInstructionCreator paymentInstructionCreator;
private final ImdgId idProvider;
private final KafkaSender kafka;
private final LiabilitiesClaimsAssetsCreator lbltsClmsAssetsCreator;
private final LiabilitiesClaimsMoneyCreator lbltsClmsMoneyCreator;
private Long lastAccBalanceResponseId = 0L;
//при прохождении по выгруженным ExecutionDeposit, ошибочные статусы проставляются для
//контр сделок. В таком случае, в коллекции хранятся не синхронизированные с IMDG ExecutionDeposit
//для которых статус должен быть DENIED
//************ !!! NOT THREAD SAFE !!! ************
private final Set<Long> deniedIds = new HashSet<>();
private Long clearingSessionId;
@Autowired
public Clearing(ImdgProvider imdgProvider,
ExecutionDepositSorter sorter, PaymentInstructionCreator paymentInstructionCreator,
BiFunction<ClearingCategory, ExecutionDeposit, IValidator> validation, KafkaSender kafka, LiabilitiesClaimsAssetsCreator liabilitiesClaimsAssetsCreator, LiabilitiesClaimsMoneyCreator lbltsClmsMoneyCreator) {
this.executionDepositImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionDeposit, ExecutionDeposit.class);
this.clearingCategoryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ClearingMemberCategory, ClearingMemberCategory.class);
this.liabilitiesClaimsAssetsImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_LiabilitiesClaimsAssets, LiabilitiesClaimsAssets.class);
this.liabilitiesClaimsMoneyImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_LiabilitiesClaimsMoney, LiabilitiesClaimsMoney.class);
this.paymentInstructionCreator = paymentInstructionCreator;
this.idProvider = imdgProvider.getImdgIdGenerator();
this.sorter = sorter;
this.validation = validation;
this.kafka = kafka;
this.lbltsClmsAssetsCreator = liabilitiesClaimsAssetsCreator;
this.lbltsClmsMoneyCreator = lbltsClmsMoneyCreator;
}
public void startClearing() {
clearingSessionId = idProvider.nextId();
log.info("startClearing clearingSessionId = {}. loading ExecutionDeposits", clearingSessionId);
//выгружаем ExecutionDeposit с пустым sessionId
Collection<ExecutionDeposit> execDeposits = executionDepositImdg.getCollectionObjectsBySQL("sessionId = null");
log.info("loaded {} ExecutionDeposits", execDeposits.size());
//делаем группировку ExecutionDeposit по категории и компании, делаем сортировку
LinkedHashMap<ClearingCategory, List<List<ExecutionDeposit>>> executionDeposits = sorter.sortExecutionDeposit(execDeposits);
executionDeposits.forEach((category, companies) -> {
for (List<ExecutionDeposit> companyExecDeposits : companies) {
runChecksSetStatus(category, companyExecDeposits);
}
});
executionDeposits.forEach((category, companies) -> {
for (List<ExecutionDeposit> companyExecDeposits : companies) {
processAllowedExecutionDeposit(category, companyExecDeposits);
}
});
}
//кладем в порядке I инициатор, V внутренний, B банк (ответная категория)
//получается I -> B; V -> B
//
//1) ТОЛЬКО ДЛЯ V: если DENIED по любым причинам то для всех следующих сделок этой компании DENIED
//2) DENIED сделка всегда встречная сделка тоже делается DENIED
//3) DENIED может быть для сделки которая уже прошла обработку и стала ALLOWED
private void runChecksSetStatus(ClearingCategory category, List<ExecutionDeposit> executionDeposits) {
log.info("setting statuses for category: {}, company: {}", category, executionDeposits.get(0).getCompanyId());
boolean financialError = false;
for (ExecutionDeposit execDeposit : executionDeposits) {
if (deniedIds.contains(execDeposit.getId())) {
log.trace("executionDeposit {} was already denied", execDeposit.getId());
continue;
}
if (financialError) {
log.trace("executionDeposit {} FinancialObligationNotSatisfied was present earlier for this company; setting DENIED status", execDeposit.getId());
updateDenied(execDeposit);
continue;
}
IValidator validator = validation.apply(category, execDeposit);
Optional<EnumMessage> errorMsg = validator.tillFirstError();
if (errorMsg.isPresent() && !ClearingErrorInternal.FinancialObligationNotSatisfied.equals(errorMsg.get().getSubject())) {
log.trace("executionDeposit {} setting DENIED status: {}", execDeposit.getId(), errorMsg.get().getSubject().nameOrId());
updateDenied(execDeposit);
continue;
} else if (errorMsg.isPresent()) {
log.trace("executionDeposit {} setting DENIED status: {}", execDeposit.getId(), errorMsg.get().getSubject().nameOrId());
financialError = true;
updateDenied(execDeposit);
continue;
}
updateAllowed(execDeposit);
log.trace("executionDeposit {} was updated ALLOWED", execDeposit.getId());
}
}
private void updateDenied(ExecutionDeposit execDeposit) {
execDeposit.setCoverageStatus(Allowed.DENIED.getKey());
execDeposit.setUpdated(Instant.now());
executionDepositImdg.update(execDeposit);
log.trace("executionDeposit {} was updated DENIED", execDeposit.getId());
//также необходимо установить статус DEND встречной сделке контрагента этой компании, которая выбирается из executionDeposit по ключу:
//securityId И чтобы сделка была компании категории clearingMemberCategory.clearingMemberCategory ≠ той категории, сделка которой обрабатывается в настоящий момент.
ExecutionDeposit matchedExecDeposit = executionDepositImdg.getSingleObjectBySQL(
"securityId = " + execDeposit.getSecurityId()
+ " AND companyId != " + execDeposit.getCompanyId());
if (matchedExecDeposit != null) {
matchedExecDeposit.setCoverageStatus(Allowed.DENIED.getKey());
matchedExecDeposit.setUpdated(Instant.now());
executionDepositImdg.update(matchedExecDeposit);
log.trace("executionDeposit {} matching executionDeposit {} was updated DENIED", execDeposit.getId(), matchedExecDeposit.getId());
deniedIds.add(matchedExecDeposit.getId());
}
}
private void updateAllowed(ExecutionDeposit execDeposit) {
execDeposit.setCoverageStatus(Allowed.ALLOWED.getKey());
execDeposit.setUpdated(Instant.now());
executionDepositImdg.update(execDeposit);
}
private void processAllowedExecutionDeposit(ClearingCategory category, List<ExecutionDeposit> executionDeposits) {
log.info("running clearing for category: {}, company: {}", category, executionDeposits.get(0).getCompanyId());
for (ExecutionDeposit execDeposit : executionDeposits) {
if (deniedIds.contains(execDeposit.getId())) {
log.trace("Skipping processing executionDeposit {} was denied as matching", execDeposit.getId());
continue;
}
if (!Allowed.ALLOWED.getKey().equals(execDeposit.getCoverageStatus())) {
log.trace("Skipping processing executionDeposit {} was denied", execDeposit.getId());
continue;
}
//send request to kafka, wait for a reply
synchronized (this) {
AccountBalanceClearingRequest request = new AccountBalanceClearingRequest();
request.setAccountId(execDeposit.getAccountId());
request.setCompanyId(execDeposit.getCompanyId());
request.setFirstLegAmount(execDeposit.getFirstLegAmount());
Long sentRequestId = kafka.sendRequestToQueue(Consts.BALANCE_ACCOUNT_UPDATE, request);
log.trace("sent request to kafka, requestId: {}", sentRequestId);
try {
wait(10000);
//если на этом месте произошла ошибка - непонятно как обрабатывать
//т.к. этом может быть уже вторая часть сделки, для первой accountBalance уже был обновлен
//либо вне зависимости от того какая часть сделки, accountBalance на самом деле мог быть
//обновлен, и была проблема была с сетью/недоступностью кафки etc.
} catch (InterruptedException e) {
log.error(ExceptionUtils.getStackTrace(e));
}
if (!Objects.equals(lastAccBalanceResponseId, sentRequestId)) {
log.error("FATAL: last response from AccountBalance update from Kafka ID was {}, but sent ID was {}",
lastAccBalanceResponseId,
sentRequestId);
continue;
}
}
LiabilitiesClaimsAssets lca = lbltsClmsAssetsCreator.createLiabilitiesClaimsAssets(category, execDeposit);
liabilitiesClaimsAssetsImdg.insert(lca);
log.trace("LiabilitiesClaimsAssets {} was created", lca.getId());
Optional<LiabilitiesClaimsMoney> lcmFirstLegFound = lbltsClmsMoneyCreator.searchLcmBySettlementDate(lca.getAccountId(), lca.getCompanyId(), lca.getSettlementDate());
var wrapperFirstLegCreated = new Object() {LiabilitiesClaimsMoney firstLegCreated;};
lcmFirstLegFound.ifPresentOrElse(lcm -> {
log.trace("LiabilitiesClaimsMoney first leg {} was found", lcm.getId());
lbltsClmsMoneyCreator.updateFirstLegLcm(lcm, lca, category);
liabilitiesClaimsMoneyImdg.update(lcm);
}, () -> {
LiabilitiesClaimsMoney lcmFirstLeg = lbltsClmsMoneyCreator.createFirstLegLcm(lca, category);
liabilitiesClaimsMoneyImdg.insert(lcmFirstLeg);
log.trace("LiabilitiesClaimsMoney first leg {} was created", lcmFirstLeg.getId());
wrapperFirstLegCreated.firstLegCreated = lcmFirstLeg;
});
Optional<LiabilitiesClaimsMoney> lcmSecondLegFound = lbltsClmsMoneyCreator.searchLcmByRefundDate(lca.getAccountId(), lca.getCompanyId(), lca.getRefundDate());
lcmSecondLegFound.ifPresentOrElse(lcm -> {
log.trace("LiabilitiesClaimsMoney second leg {} was found", lcm.getId());
lbltsClmsMoneyCreator.updateSecondLegLcm(lcm, lca, category);
liabilitiesClaimsMoneyImdg.update(lcm);
}, () -> {
LiabilitiesClaimsMoney lcmSecondLeg = lbltsClmsMoneyCreator.createSecondLegLcm(lca, category);
liabilitiesClaimsMoneyImdg.insert(lcmSecondLeg);
log.trace("LiabilitiesClaimsMoney second leg {} was created", lcmSecondLeg.getId());
});
lbltsClmsAssetsCreator.updateLca(lca, lcmFirstLegFound.orElse(wrapperFirstLegCreated.firstLegCreated));
liabilitiesClaimsAssetsImdg.update(lca);
log.trace("creating paymentInstructions");
List<PaymentInstruction> paymentInstructions = paymentInstructionCreator.createAndSavePaymentInstructions(lca, Collections.singleton(category));
paymentInstructionCreator.updateFieldLiabilitiesClaimsAssetsByPayment(lca, paymentInstructions, Collections.singleton(category));
liabilitiesClaimsAssetsImdg.update(lca);
var firstLegLiabilitiesClaimsMoney = lcmFirstLegFound.orElseGet(() -> wrapperFirstLegCreated.firstLegCreated);
log.trace("updating LiabilitiesClaimsMoney by payment");
paymentInstructionCreator.updateFieldLiabilitiesClaimsMoneyByPayment(firstLegLiabilitiesClaimsMoney, paymentInstructions, Collections.singleton(category));
liabilitiesClaimsMoneyImdg.update(firstLegLiabilitiesClaimsMoney);
log.trace("LiabilitiesClaimsMoney first leg {} was updated", firstLegLiabilitiesClaimsMoney.getId());
}
}
public void setLastAccBalanceResponseId(Long lastAccBalanceResponseId) {
this.lastAccBalanceResponseId = lastAccBalanceResponseId;
}
}

View file

@ -5,8 +5,9 @@ import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.DisposableBean;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.scheduling.annotation.EnableScheduling;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Service;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonIdRequest;
import ru.spcex.platform.utils.log.ExceptionUtils;
import java.util.concurrent.ExecutorService;
@ -20,13 +21,18 @@ public class ClearingService implements DisposableBean {
private final SdfCreatorBySTLDPayment sdfCreator;
private final PaymentUpdateBySdf04 paymentUpdater;
private final VerificationResultComponent verificationResultComponent;
private final Clearing clearing;
private final ExecutionDepositComponent executionDepositComponent;
@Autowired
public ClearingService(SdfCreatorBySTLDPayment sdfCreator, PaymentUpdateBySdf04 paymentUpdater,
VerificationResultComponent verificationResultComponent) {
VerificationResultComponent verificationResultComponent, Clearing clearing,
ExecutionDepositComponent executionDepositComponent) {
this.sdfCreator = sdfCreator;
this.paymentUpdater = paymentUpdater;
this.verificationResultComponent = verificationResultComponent;
this.executionDepositComponent = executionDepositComponent;
this.clearing = clearing;
this.executor = Executors.newSingleThreadExecutor();
}
@ -69,4 +75,35 @@ public class ClearingService implements DisposableBean {
log.debug("Shutdown {}", getClass().getSimpleName());
executor.shutdown();
}
public void executeClearing() {
log.info("start clearing");
executor.execute(() -> {
try {
clearing.startClearing();
} catch (Throwable e) {
log.error("{}", ExceptionUtils.getStackTrace(e));
}
});
}
public void continueClearing(BaseRequest<CommonIdRequest> event) {
synchronized (clearing) {
CommonIdRequest requestPayload = event.getRequestPayload();
clearing.setLastAccBalanceResponseId(requestPayload.getId());
clearing.notifyAll();
}
}
public void executeSTrade() {
log.info("start STrade check for ExecutionDeposit");
executor.execute(() -> {
try {
executionDepositComponent.processNewTS();
} catch (Throwable e) {
log.error("{}", ExceptionUtils.getStackTrace(e));
}
});
}
}

View file

@ -5,6 +5,7 @@ import org.springframework.beans.factory.InitializingBean;
import org.springframework.stereotype.Service;
import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.Sdf04Request;
import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonIdRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
import ru.spcex.platform.enumeration.Task;
@ -31,6 +32,15 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
callback(LauncherCommandRequest.class)
.setConsumer(event -> clearingService.executeVerification())
.forDestination(Task.getVerification.topic(), callbacks::put);
callback(Object.class) //todo check Object suitable
.setConsumer(event -> clearingService.executeClearing())
.forDestination(Task.startOfClearing.topic(), callbacks::put);
callback(CommonIdRequest.class)
.setConsumer(clearingService::continueClearing)
.forDestination(Consts.CONTINUE_CLEARING, callbacks::put);
callback(Object.class)
.setConsumer(event -> clearingService.executeSTrade())
.forDestination(Task.getOfTrades.topic(), callbacks::put);
init();
}
}

View file

@ -0,0 +1,26 @@
package ru.spcex.clearing.service;
import ru.clearing.classes.statics.data.execution.ExecutionDeposit;
import ru.spcex.platform.enumeration.ClearingCategory;
public class ExecDepositWCategory {
private final ExecutionDeposit executionDeposit;
private final ClearingCategory clearingMemberCategory;
public ExecDepositWCategory(ExecutionDeposit executionDeposit, ClearingCategory clearingMemberCategory) {
this.executionDeposit = executionDeposit;
this.clearingMemberCategory = clearingMemberCategory;
}
public ExecutionDeposit getExecutionDeposit() {
return executionDeposit;
}
public ClearingCategory getClearingMemberCategory() {
return clearingMemberCategory;
}
public Long getCompanyId() {
return executionDeposit.getCompanyId();
}
}

View file

@ -1,6 +1,5 @@
package ru.spcex.clearing.service;
import com.fasterxml.jackson.annotation.JsonProperty;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.clients.producer.RecordMetadata;
@ -10,15 +9,11 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.scheduling.annotation.EnableScheduling;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
import org.springframework.util.StringUtils;
import ru.clearing.classes.statics.data.account.Account;
import ru.clearing.classes.statics.data.account.AccountBalance;
import ru.clearing.classes.statics.data.clearing.VerificationResult;
import ru.clearing.classes.statics.data.company.Company;
import ru.clearing.classes.statics.data.execution.ExecutionDeposit;
import ru.clearing.classes.statics.data.misc.Listing;
import ru.clearing.classes.statics.data.misc.STrade;
import ru.clearing.classes.statics.data.sdf.SDf01;
import ru.clearing.classes.statics.data.security.Security;
import ru.spcex.clearing.error.ClearingError;
import ru.spcex.clearing.error.ClearingException;
@ -26,12 +21,11 @@ import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.ActionType;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.registry.DealRegisterNewRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.securitites.MoneyMarketSecurityNewRequest;
import ru.spcex.platform.classes.base.interfaces.WithId;
import ru.spcex.platform.enumeration.AccountType;
import ru.spcex.platform.classes.base.SpcexObjectBase;
import ru.spcex.platform.enumeration.Allowed;
import ru.spcex.platform.enumeration.Market;
import ru.spcex.platform.enumeration.ResultStatuses;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgId;
import ru.spcex.platform.imdg.api.ImdgProvider;
@ -40,17 +34,13 @@ import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.log.ExceptionUtils;
import ru.spcex.platform.utils.time.TimeUtil;
import ru.spcex.clearing.platform.messaging.domain.cud.registry.CoveredDealRegisterNewRequest;
import java.math.BigDecimal;
import java.math.RoundingMode;
import java.time.Instant;
import java.time.LocalDate;
import java.util.*;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.Future;
import java.util.function.Function;
import java.util.stream.Collectors;
/**
@ -96,8 +86,9 @@ public class ExecutionDepositComponent {
/**
* Сбрасывать каждый день в 01:00:01 "0 1 0 1 * ?"
*/
@Scheduled(cron = "${clearing-service.scheduler.check-s-trade}")
@Scheduled(cron = "0 1 0 1 * ?")
public void resetTradingDay() {
log.trace("Recheck today trading day for search STrade. Current state: tradeNum={}, tradeDat={}", tradeNum, tradingDay);
Instant today = TimeUtil.localDateToInstant(LocalDate.now());
if (tradingDay == null || !tradingDay.equals(today)) {
tradeNum = -1L;
@ -107,9 +98,9 @@ public class ExecutionDepositComponent {
}
public void processNewTS() {
Long tradeNum = -1L; // todo уточнить как он обновляется
log.debug("Start check new S_TRADE after {}", tradingDay);
ImdgPredicateBuilder pb = sTradeImdg.predicateBuilder();
ImdgPredicate sql = pb.and(pb.greater("tradeNum",tradeNum), pb.greatEqual("tradeDateTime", tradingDay));
ImdgPredicate sql = pb.greatEqual("tradeDateTime", tradingDay);
Collection<STrade> sTrades = sTradeImdg.getCollectionObjectsByPredicate(sql);
log.info("Found {} new s_trade with trade_num>{}", sTrades.size(), tradeNum);
@ -121,12 +112,13 @@ public class ExecutionDepositComponent {
// Выявление новых сделок необходимо выполнить следующие контрольные проверки:
// Проверить все инструменты.
Set<String> secCodesOfSecurity;
{
Set<String> secCodesOfSTrade = sTrades.stream().map(STrade::getSecCode).filter(Objects::nonNull).collect(Collectors.toSet());
log.debug("Verify {} instruments for {} STrade's.", secCodesOfSTrade.size(), sTrades.size());
ImdgPredicate allIn = securityImdg.predicateBuilder().in("securitySymbol", secCodesOfSTrade.toArray(new String[0]));
Collection<Security> foundSecurities = securityImdg.getCollectionObjectsByPredicate(allIn);
Set<String> secCodesOfSecurity = foundSecurities.stream().map(Security::getSecuritySymbol).filter(Objects::nonNull).collect(Collectors.toSet());
secCodesOfSecurity = foundSecurities.stream().map(Security::getSecuritySymbol).filter(Objects::nonNull).collect(Collectors.toSet());
if (secCodesOfSecurity.containsAll(secCodesOfSTrade)) {
log.debug("All {} Security found by {} secCodes from STrade",
secCodesOfSecurity.size(), secCodesOfSTrade.size());
@ -135,59 +127,49 @@ public class ExecutionDepositComponent {
notFoundSymbol.removeAll(secCodesOfSecurity);
log.info("Found only {} Security by {} secCodes from STrade. Not found: {}",
secCodesOfSecurity.size(), secCodesOfSTrade.size(), notFoundSymbol);
createNewSecurities(notFoundSymbol);
log.info("Stop till they all will be created");
auditMessage("В security нет записей с securitySymbol", notFoundSymbol);
return;
}
}
Long generationId = idGenerator.nextId();
log.info("generationId = {}", generationId);
for (STrade trade:sTrades) {
log.trace("Check s_trade[{}].tradeNum={}", trade.getId(), trade.getTradeNum());
LocalDate today = LocalDate.now();
for (STrade trade : sTrades) {
log.trace("Check s_trade[{}].tradeNum={} on date {}", trade.getId(), trade.getTradeNum(), today);
Collection<ExecutionDeposit> existsEDeposit = executionDepositImdg.getCollectionObjectsByFieldValues(Map.of(
"exchangeExecutionId", trade.getTradeNum(),
"exchangeExecutionTime", trade.getTradeDateTime()
"clearingDate", today
));
if (existsEDeposit.isEmpty()) {
log.trace("S_TRADE[{}] new", trade.getId());
ExecutionDeposit newED = null;
// Проверка secCode
if (!secCodesOfSecurity.contains(trade.getSecCode())) {
log.error("Error {}: STrade[{}].secCode={} not found",
ClearingError.RecordNotFound.getId(), trade.getId(), trade.getSecCode());
continue;
}
ExecutionDeposit newED;
try {
newED = createExecutionDeposit(trade, Allowed.ALLOWED/*todo уточнить момент заполнения*/, generationId);
verification(newED);
newED = createExecutionDeposit(trade, null, null);
executionDepositImdg.insert(newED);
sendNotification(newED);
} catch (ClearingException ce) {
auditMessage(ce);
} catch (Exception e) {
if (newED != null) {
newED.setCoverageStatus(Allowed.DENIED.getKey());
}
log.error("When create new ExecutionDeposit by STrade[{}]", trade.getId());
log.error("When create new ExecutionDeposit by STrade[{}] error: {}", trade.getId(), ExceptionUtils.getStackTrace(e));
}
} else {
long[] idToLong = existsEDeposit.stream().mapToLong(ed-> ed.getId()).toArray();
log.warn("S_TRADE[{}] already has executionDeposit: {}", trade.getId(), Arrays.toString(idToLong));
long[] idToLong = existsEDeposit.stream().mapToLong(SpcexObjectBase::getId).toArray();
log.trace("S_TRADE[{}] already has executionDeposit: {}", trade.getId(), Arrays.toString(idToLong));
}
}
Long newMaxTradeNum = sTrades.stream().mapToLong(STrade::getTradeNum).max().orElseGet(()-> tradeNum);
log.debug("Next tradeNum is {}", newMaxTradeNum);
Long newMaxTradeNum = sTrades.stream().mapToLong(STrade::getTradeNum).max().orElseGet(() -> tradeNum);
log.info("Process completed. Next tradeNum is {}", newMaxTradeNum);
}
protected void verification(ExecutionDeposit forED) throws ClearingException {
/*todo Рассчитанные в КС контрольные суммы (общее количество сделок и суммарный объем заключенных сделок в денежном выражении)
должны совпадать со значениями, рассчитанными Торговой системой:
count(execution[tradingDay]) = count (trade_arqua)
*/
// использовать ли VerificationResultComponent для сверки или здесь код добавить.
}
protected void auditMessage(ClearingException ce) {
log.error("AUDIT error code {}: {}", ce.getEnumMsg(), ce.getMessage());
@ -201,71 +183,12 @@ public class ExecutionDepositComponent {
log.error("audit \"clearing-service\", errorText: {}", txt);
}
/**
* в очередь kafka для модуля securities-service сообщение о добавлении инструмента с параметром securitySymbol=s_trade.sec_code
* @param newSymbolRequest
*/
protected void createNewSecurities(Collection<String> newSymbolRequest) {
final String destination = Consts.DESTINATION_MONEY_MARKET_SECURITY_NEW;
List<String> symbolRequests = new ArrayList<>(newSymbolRequest); // чтобы в случае ошибки отобразить номер в логе
List<Future<RecordMetadata>> sendAll = new ArrayList<>(symbolRequests.size());
for (String newSymbol: symbolRequests) {
if (newSymbol == null || newSymbol.isEmpty()) {
log.warn("Empty SecuritySumbol");
} else {
MoneyMarketSecurityNewRequest requestPayload = new MoneyMarketSecurityNewRequest();
requestPayload.setSecuritySymbol(newSymbol);
BaseRequest<Object> request = new BaseRequest<>();
request.setId(idGenerator.nextId());
request.setActionType(ActionType.NEW);
request.setRequestPayload(requestPayload);
// saveRequestToStorage(destination, request); //сохраняет данные о запросе в хранилище
log.trace("Send to {} new symbol \"{}\" ", destination, newSymbol);
Future<RecordMetadata> send = kafka.send(new ProducerRecord<>(destination, request));
sendAll.add(send);
}
}
int i = 0;
for (Future<RecordMetadata> future: sendAll) {
try {
future.get(); // get exception
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
String about = i < symbolRequests.size() ? symbolRequests.get(i) : "(out of range i=" + i + ")";
log.warn("Thread interrupted! On send symbol \"{}\"", about);
throw new RuntimeException(e);
} catch (ExecutionException e) {
String about = i < symbolRequests.size() ? symbolRequests.get(i) : "(out of range i=" + i + ")";
log.error("Error send message for symbol \"{}\" to {}: {}", destination,
about, ExceptionUtils.getStackTrace(e.getCause() == null ? e : e.getCause()));
}
i++;
}
}
protected void sendNotification(ExecutionDeposit forED) {
final String destination = Consts.REGISTRY_COVERED_DEAL_REGISTER_NEW;
CoveredDealRegisterNewRequest requestPayload = new CoveredDealRegisterNewRequest();
final String destination = Consts.REGISTRY_DEAL_REGISTER_NEW;
DealRegisterNewRequest requestPayload = new DealRegisterNewRequest();
requestPayload.setExecutionId(forED.getId());
// requestPayload.setCompanyFullName(forED.getCompanyFullName());
requestPayload.setTradingDate(forED.getTradingDate());
requestPayload.setExchangeExecutionId(forED.getExchangeExecutionId());
requestPayload.setExchangeExecutionTime(forED.getExchangeExecutionTime());
requestPayload.setSecuritySymbol(forED.getSecuritySymbol());
requestPayload.setSecurityFullName(forED.getSecurityFullName());
// requestPayload.setSellerFullName(forED.getSellerFullName());
// requestPayload.setSellerClearingCode(forED.getSellerClearingCode());
// String requestPayload.setSellerAccount(forED.getAccountId());
// String requestPayload.setBuyerFullName(forED.getBuyerFullName());
// requestPayload.setBuyerClearingCode(forED.getBuyerClearingCode());
// String requestPayload.setBuyerAccount(forED.getBuyerAccount());
// BigDecimal requestPayload.setAmount(forED.getAmount());
requestPayload.setId(idGenerator.nextId());
requestPayload.setCreatedAt(forED.getCreated());
requestPayload.setUpdatedAt(forED.getUpdated());
requestPayload.setClearingDate(forED.getClearingDate());
// requestPayload.setId(idGenerator.nextId());
BaseRequest<Object> request = new BaseRequest<>();
request.setId(idGenerator.nextId());

View file

@ -0,0 +1,33 @@
package ru.spcex.clearing.service;
import org.apache.kafka.clients.consumer.Consumer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.InitializingBean;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
import ru.spcex.platform.enumeration.Task;
@Service
public class LauncherCommandReceiver extends QueueConsumer implements InitializingBean {
private final Logger log = LoggerFactory.getLogger(getClass());
private final ClearingService clearingService;
@Autowired
public LauncherCommandReceiver(Consumer<String, Object> kafkaQueue,
ClearingService clearingService) {
super(kafkaQueue);
this.clearingService = clearingService;
}
@Override
public void afterPropertiesSet() {
callback(LauncherCommandRequest.class)
.setConsumer(action -> clearingService.executeVerification())
.forDestination(Task.createOrderConfirm.topic(), callbacks::put); // CORC
init();
}
}

View file

@ -11,10 +11,10 @@ import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.SdfClearingRequest;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.clearing.service.builder.Sdf03Builder;
import ru.spcex.clearing.service.builder.Sdf11Builder;
import ru.spcex.clearing.service.order.PaymentBatchInfo;
import ru.spcex.clearing.service.order.PaymentInstructionSorter;
import ru.spcex.clearing.service.order.Sdf03Builder;
import ru.spcex.clearing.service.order.Sdf11Builder;
import ru.spcex.platform.enumeration.TransactionStatus;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgId;

View file

@ -0,0 +1,76 @@
package ru.spcex.clearing.service.builder;
import org.springframework.beans.factory.annotation.Autowired;
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.liabilities.LiabilitiesClaimsAssets;
import ru.clearing.classes.statics.data.liabilities.LiabilitiesClaimsMoney;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.platform.enumeration.AccountType;
import ru.spcex.platform.enumeration.ClearingCategory;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.utils.enumeration.IEnumKey;
import java.time.Instant;
import java.time.LocalDate;
@Component
public class LiabilitiesClaimsAssetsCreator {
private final Imdg<Account> accountImdg;
private final Imdg<Company> companyImdg;
@Autowired
public LiabilitiesClaimsAssetsCreator(ImdgProvider imdgProvider) {
this.accountImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Account, Account.class);
this.companyImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Company, Company.class);
}
public LiabilitiesClaimsAssets createLiabilitiesClaimsAssets(ClearingCategory category, ExecutionDeposit executionDeposit) {
LiabilitiesClaimsAssets liabilitiesClaimsAssets = new LiabilitiesClaimsAssets();
liabilitiesClaimsAssets.setCompanyId(executionDeposit.getCompanyId());
liabilitiesClaimsAssets.setAccountId(executionDeposit.getAccountId());
Account account = accountImdg.getSingleObjectByID(executionDeposit.getAccountId());
liabilitiesClaimsAssets.setAccountType(account.getAccountType());
liabilitiesClaimsAssets.setAccount(account.getAccount());
AccountType accType = IEnumKey.getEnumByKey(AccountType.class, account.getAccountType());
boolean iClrn = ClearingCategory.I.equals(category) && AccountType.Clrn.equals(accType);
boolean vInfo = ClearingCategory.V.equals(category) && AccountType.Info.equals(accType);
boolean bClrn = ClearingCategory.B.equals(category) && AccountType.Clrn.equals(accType);
if (iClrn || vInfo) {
liabilitiesClaimsAssets.setLiabilitiesQuantity(executionDeposit.getFirstLegAmount());
}
if (bClrn) {
liabilitiesClaimsAssets.setClaimsQuantity(executionDeposit.getSecondLegAmount());
}
liabilitiesClaimsAssets.setCurrency(executionDeposit.getSettlementCurrency());
liabilitiesClaimsAssets.setSettlementDate(executionDeposit.getFirstLegSettlementDate());
liabilitiesClaimsAssets.setTradingDate(executionDeposit.getTradingDate());
liabilitiesClaimsAssets.setRefundDate(executionDeposit.getSecondLegSettlementDate());
liabilitiesClaimsAssets.setPrice(executionDeposit.getPrice());
liabilitiesClaimsAssets.setSecurityId(executionDeposit.getSecurityId());
Company company = companyImdg.getSingleObjectByID(executionDeposit.getCompanyId());
liabilitiesClaimsAssets.setTradingCode(company.getTradingCode());
liabilitiesClaimsAssets.setClearingCode(company.getClearingCode());
liabilitiesClaimsAssets.setShortName(company.getShortName());
liabilitiesClaimsAssets.setContract(executionDeposit.getSecuritySymbol());
liabilitiesClaimsAssets.setFullName(company.getFullName());
//todo liabilitiesClaimsAssets.setLiabilitiesClaimsMoneyId(company.getFullName());
//fixme liabilitiesClaimsAssets.setClearingStatus();
//todo liabilitiesClaimsAssets.setPaymentId();
//todo liabilitiesClaimsAssets.setRefundPaymentId();
liabilitiesClaimsAssets.setCreated(Instant.now());
liabilitiesClaimsAssets.setClearingDate(LocalDate.now());
return liabilitiesClaimsAssets;
}
public void updateLca(LiabilitiesClaimsAssets lca,
LiabilitiesClaimsMoney firstLeg) {
lca.setLiabilitiesClaimsMoneyId(firstLeg.getId());
lca.setUpdated(Instant.now());
}
}

View file

@ -0,0 +1,104 @@
package ru.spcex.clearing.service.builder;
import org.springframework.stereotype.Component;
import ru.clearing.classes.statics.data.liabilities.LiabilitiesClaimsAssets;
import ru.clearing.classes.statics.data.liabilities.LiabilitiesClaimsMoney;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.platform.enumeration.ClearingCategory;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import java.math.BigDecimal;
import java.time.LocalDate;
import java.time.format.DateTimeFormatter;
import java.util.Optional;
import static ru.spcex.platform.utils.number.BigDecimalUtil.safeBD;
@Component
public class LiabilitiesClaimsMoneyCreator {
private final Imdg<LiabilitiesClaimsMoney> liabilitiesClaimsMoneyImdg;
private final static DateTimeFormatter formatter = DateTimeFormatter.ofPattern("yyyy-MM-dd");
public LiabilitiesClaimsMoneyCreator(ImdgProvider imdgProvider) {
this.liabilitiesClaimsMoneyImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_LiabilitiesClaimsMoney, LiabilitiesClaimsMoney.class);
}
public Optional<LiabilitiesClaimsMoney> searchLcmBySettlementDate(Long accountId, Long companyId, LocalDate settlementDate) {
String sql = String.format("accountId=%d and companyId=%d and settlementDate=%s",
accountId, companyId, formatter.format(settlementDate));
LiabilitiesClaimsMoney liabilitiesClaimsMoney = liabilitiesClaimsMoneyImdg.getSingleObjectBySQL(sql);
return Optional.ofNullable(liabilitiesClaimsMoney);
}
public void updateFirstLegLcm(LiabilitiesClaimsMoney liabilitiesClaimsMoney,
LiabilitiesClaimsAssets liabilitiesClaimsAssets, ClearingCategory category) {
BigDecimal prevLAmount = safeBD(liabilitiesClaimsMoney.getLiabilitiesAmount());
BigDecimal prevClaimsAmount = safeBD(liabilitiesClaimsMoney.getClaimsAmount());
BigDecimal lcaLiabilitiesQuantity = safeBD(liabilitiesClaimsAssets.getLiabilitiesQuantity());
BigDecimal lcaClaimsQuantity = safeBD(liabilitiesClaimsAssets.getClaimsQuantity());
liabilitiesClaimsMoney.setLiabilitiesAmount(prevLAmount.add(lcaLiabilitiesQuantity));
liabilitiesClaimsMoney.setClaimsAmount(prevClaimsAmount.add(lcaClaimsQuantity));
}
public LiabilitiesClaimsMoney createFirstLegLcm(LiabilitiesClaimsAssets lca, ClearingCategory category) {
LiabilitiesClaimsMoney lcm = new LiabilitiesClaimsMoney();
lcm.setCompanyId(lca.getCompanyId());
lcm.setAccountId(lca.getAccountId());
lcm.setAccountType(lca.getAccountType());
lcm.setAccount(lca.getAccount());
if (category.equals(ClearingCategory.I) || category.equals(ClearingCategory.V)) {
lcm.setLiabilitiesAmount(lca.getLiabilitiesQuantity());
} else if (category.equals(ClearingCategory.B)) {
lcm.setClaimsAmount(lca.getClaimsQuantity());
}
lcm.setSettlementDate(lca.getSettlementDate());
lcm.setTradingDate(lca.getTradingDate());
lcm.setCurrency(lca.getCurrency());
lcm.setTradingCode(lca.getTradingCode());
lcm.setShortName(lca.getShortName());
lcm.setFullName(lca.getFullName());
lcm.setCreated(lca.getCreated());
lcm.setClearingDate(LocalDate.now());
return lcm;
}
public Optional<LiabilitiesClaimsMoney> searchLcmByRefundDate(Long accountId, Long companyId, LocalDate refundDate) {
String sql = String.format("accountId=%d and companyId=%d and settlementDate=%s",
accountId, companyId, formatter.format(refundDate));
LiabilitiesClaimsMoney liabilitiesClaimsMoney = liabilitiesClaimsMoneyImdg.getSingleObjectBySQL(sql);
return Optional.ofNullable(liabilitiesClaimsMoney);
}
public LiabilitiesClaimsMoney createSecondLegLcm(LiabilitiesClaimsAssets lca, ClearingCategory category) {
LiabilitiesClaimsMoney lcm = new LiabilitiesClaimsMoney();
lcm.setCompanyId(lca.getCompanyId());
lcm.setAccountId(lca.getAccountId());
lcm.setAccountType(lca.getAccountType());
lcm.setAccount(lca.getAccount());
if (category.equals(ClearingCategory.I) || category.equals(ClearingCategory.V)) {
lcm.setClaimsAmount(lca.getClaimsQuantity());
} else if (category.equals(ClearingCategory.B)) {
lcm.setLiabilitiesAmount(lca.getLiabilitiesQuantity());
}
lcm.setSettlementDate(lca.getRefundDate());
lcm.setTradingDate(lca.getTradingDate());
lcm.setCurrency(lca.getCurrency());
lcm.setTradingCode(lca.getTradingCode());
lcm.setShortName(lca.getShortName());
lcm.setFullName(lca.getFullName());
lcm.setCreated(lca.getCreated());
lcm.setClearingDate(lca.getRefundDate());
return lcm;
}
public void updateSecondLegLcm(LiabilitiesClaimsMoney lcm, LiabilitiesClaimsAssets lca, ClearingCategory category) {
BigDecimal prevLiabilitiesAmount = safeBD(lcm.getLiabilitiesAmount());
BigDecimal prevClaimsAmount = safeBD(lcm.getClaimsAmount());
BigDecimal lcaLiabilitiesQuantity = safeBD(lca.getLiabilitiesQuantity());
BigDecimal lcaClaimsQuantity = safeBD(lca.getClaimsQuantity());
lcm.setLiabilitiesAmount(prevLiabilitiesAmount.add(lcaLiabilitiesQuantity));
lcm.setClaimsAmount(prevClaimsAmount.add(lcaClaimsQuantity));
}
}

View file

@ -0,0 +1,533 @@
package ru.spcex.clearing.service.builder;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
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.liabilities.LiabilitiesClaimsAssets;
import ru.clearing.classes.statics.data.liabilities.LiabilitiesClaimsMoney;
import ru.clearing.classes.statics.data.payment.PaymentInstruction;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.platform.classes.base.SpcexObjectBase;
import ru.spcex.platform.enumeration.*;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgId;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.utils.time.TimeUtil;
import java.math.BigDecimal;
import java.time.Instant;
import java.time.LocalDate;
import java.time.format.DateTimeFormatter;
import java.util.*;
import java.util.concurrent.atomic.AtomicLong;
import java.util.stream.Collectors;
/*
todo использовать так:
paymentInstructionCreator.createAndSavePaymentInstructions()
paymentInstructionCreator.updateLiabilitiesClaimsAssetsByPaymentInstruction();
paymentInstructionCreator.updateLiabilitiesClaimsMoneyByPaymentInstruction();
*/
@Service
public class PaymentInstructionCreator {
protected static final Long SPVB_ID = 1L; // СПВБ
protected static final Long PRC_ID = 2L; // 2 "НКО АО ПРЦ"
private static final DateTimeFormatter DATE_FORMATTER_ddMMyy = DateTimeFormatter.ofPattern("ddMMyy");
private final Logger log = LoggerFactory.getLogger(getClass());
protected final ImdgId idGenerator;
protected final Imdg<LiabilitiesClaimsAssets> liabilitiesClaimsAssetsImdg;
protected final Imdg<CompanySymbols> companySymbolsImdg;
protected final Imdg<Company> companyImdg;
protected final Imdg<Account> accountImdg;
protected final Imdg<PaymentInstruction> paymentInstructionImdg;
protected final Imdg<LiabilitiesClaimsMoney> liabilitiesClaimsMoneyImdg;
protected AtomicLong documentNumberId = new AtomicLong(0L); // порядковый номер (сквозной по всем компаниям за день
protected LocalDate documentNumberResetAt;
@Autowired
public PaymentInstructionCreator(ImdgProvider imdgProvider) {
this.liabilitiesClaimsAssetsImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_LiabilitiesClaimsAssets, LiabilitiesClaimsAssets.class);
this.companySymbolsImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_CompanySymbols, CompanySymbols.class);
this.companyImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Company, Company.class);
this.accountImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Account, Account.class);
this.paymentInstructionImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_PaymentInstruction, PaymentInstruction.class);
this.liabilitiesClaimsMoneyImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_LiabilitiesClaimsMoney, LiabilitiesClaimsMoney.class);
this.idGenerator = imdgProvider.getImdgIdGenerator();
}
/**
* 2.2.2. Выполнить Создание paymentInstruction по liabilitiesClaimsAssets
* @return Созданные PaymentInstruction, может быть несколько. Они уже сохранены в IMDG.
*/
public List<PaymentInstruction> createAndSavePaymentInstructions(LiabilitiesClaimsAssets liabilitiesClaimsAssets,
Collection<ClearingCategory> clearingCategory) {
List<PaymentInstruction> newPaymentInstructions = createPaymentInstructions(liabilitiesClaimsAssets, clearingCategory);
log.trace("Created {} PaymentInstruction by LiabilitiesClaimsAssets[{}]. Do save ti IMDG.",
newPaymentInstructions.size(), liabilitiesClaimsAssets.getId());
for (PaymentInstruction instruction : newPaymentInstructions) {
paymentInstructionImdg.insert(instruction);
}
return newPaymentInstructions;
}
/**
* Создаёт и сохраняет
*
* @param liabilitiesClaimsAssets
* @param clearingCategory
* @return
*/
protected List<PaymentInstruction> createPaymentInstructions(LiabilitiesClaimsAssets liabilitiesClaimsAssets,
Collection<ClearingCategory> clearingCategory) {
List<PaymentInstruction> newPaymentInstructions = new ArrayList<>();
Instant now = Instant.now();
// I
if (clearingCategory.contains(ClearingCategory.I)) {
log.debug("By liabilitiesClaimsAssets[{}] create for category I 1 record", liabilitiesClaimsAssets.getId());
PaymentInstruction payment1 = new PaymentInstruction();
PaymentInstruction payment2 = new PaymentInstruction();
payment1.setId(idGenerator.nextId());
payment1.setCreated(now);
payment1.setClearingDate(TimeUtil.toLocalDate(now));
payment2.setId(idGenerator.nextId());
payment2.setCreated(now);
payment2.setClearingDate(TimeUtil.toLocalDate(now));
payment1.setSenderId(liabilitiesClaimsAssets.getCompanyId());
payment2.setSenderId(SPVB_ID); // СПВБ
payment1.setAddresseeId(SPVB_ID); // СПВБ
String sqlLCA = String.format("securityId = %s and companyId <> %s",
liabilitiesClaimsAssets.getSecurityId(), liabilitiesClaimsAssets.getCompanyId());
LiabilitiesClaimsAssets counterLiabilitiesClaimsAssets = liabilitiesClaimsAssetsImdg.getSingleObjectBySQL(sqlLCA);
if (counterLiabilitiesClaimsAssets == null) {
log.warn("Counter LiabilitiesClaimsAssets not found for : {}", sqlLCA);
} else {
payment2.setAddresseeId(counterLiabilitiesClaimsAssets.getCompanyId());
}
String symbol1 = selectSymbolValue(payment1.getAddresseeId(), CompanySymbol.BIC);
if (symbol1 == null) {
log.warn("CompanySymbols BIC not found for companyId={}", payment1.getAddresseeId());
} else {
payment1.setAdresseeBic(symbol1);
}
String symbol2 = selectSymbolValue(payment2.getAddresseeId(), CompanySymbol.BIC);
if (symbol2 == null) {
log.warn("CompanySymbols BIC not found for companyId={}", payment2.getAddresseeId());
} else {
payment2.setAdresseeBic(symbol2);
}
Company companyPRC = companyImdg.getSingleObjectByID(PRC_ID); // 2 "НКО АО ПРЦ"
String companyPRCName = null;
if (companyPRC == null) {
log.warn("Company.id=2 not found");
} else {
companyPRCName = companyPRC.getShortName();
}
payment1.setPayeeBankName(companyPRCName);
payment2.setPayeeBankName(companyPRCName);
String symbolPRC = selectSymbolValue(PRC_ID, CompanySymbol.BIC); // 2 "НКО АО ПРЦ"
payment1.setPayeeBic(symbolPRC);
payment2.setPayeeBic(symbolPRC);
payment1.setAddresseeBankName(companyPRCName);
payment2.setAddresseeBankName(companyPRCName);
payment1.setPaymentDate(TimeUtil.localDateToInstant(liabilitiesClaimsAssets.getSettlementDate()));
payment2.setPaymentDate(TimeUtil.localDateToInstant(liabilitiesClaimsAssets.getSettlementDate()));
payment1.setPaymentPurpose("Размещение депозита " + liabilitiesClaimsAssets.getContract());
payment2.setPaymentPurpose("Размещение депозита " + liabilitiesClaimsAssets.getContract());
payment1.setSettlementDate(liabilitiesClaimsAssets.getSettlementDate());
payment2.setSettlementDate(liabilitiesClaimsAssets.getSettlementDate());
Long amount = liabilitiesClaimsAssets.getLiabilitiesQuantity() == null ? null : liabilitiesClaimsAssets.getLiabilitiesQuantity().abs().longValue();
payment1.setCreditLegAmount(amount);
payment2.setCreditLegAmount(amount);
payment1.setDebitLegAmount(amount);
payment2.setDebitLegAmount(amount);
payment1.setCreditLegAccountId(liabilitiesClaimsAssets.getAccountId());
String accSql = String.format("companyId = %s AND accountType=TRAN AND accountStatus=ACTV AND processingSign=ALWD",
payment2.getSenderId());
Account anotherAcc = accountImdg.getSingleObjectBySQL(accSql);
if (anotherAcc == null) {
log.warn("Account not found: {}", accSql);
} else {
payment2.setCreditLegAccountId(anotherAcc.getId());
}
payment1.setCreditCsAccount(null);
payment2.setCreditCsAccount(null);
{
Account acc1 = accountImdg.getSingleObjectByID(payment1.getCreditLegAccountId());
if (acc1 != null) {
payment1.setCreditLegAccount(acc1.getAccount());
}
Account acc2 = accountImdg.getSingleObjectByID(payment2.getCreditLegAccountId());
if (acc2 != null) {
payment1.setCreditLegAccount(acc2.getAccount());
}
}
{
Account acc1 = selectAccount(payment1.getAddresseeId(), AccountType.Tran, Status.Active, Allowed.ALLOWED);
if (acc1 != null) {
payment1.setDebitLegAccountId(acc1.getId());
payment1.setDebitLegAccount(acc1.getAccount());
}
Account acc2 = selectAccount(payment2.getAddresseeId(), AccountType.Clrn, Status.Active, Allowed.ALLOWED);
if (acc2 != null) {
payment2.setDebitLegAccountId(acc2.getId());
payment2.setDebitLegAccount(acc2.getAccount());
}
}
payment1.setDebitCsAccount(null);
payment2.setDebitCsAccount(null);
// payment1.setCreditLegDirection(InOutDirection.out.getKey());
// payment2.setCreditLegDirection(InOutDirection.out.getKey());
// payment1.setDebitLegDirection(InOutDirection.in.getKey());
// payment2.setDebitLegDirection(InOutDirection.in.getKey());
payment1.setCreditLegCurrencyCode(liabilitiesClaimsAssets.getCurrency());
payment2.setCreditLegCurrencyCode(liabilitiesClaimsAssets.getCurrency());
payment1.setDebitLegCurrencyCode(liabilitiesClaimsAssets.getCurrency());
payment2.setDebitLegCurrencyCode(liabilitiesClaimsAssets.getCurrency());
payment1.setTransactionStatus(TransactionStatus.stld.getKey());
payment2.setTransactionStatus(TransactionStatus.stld.getKey());
payment1.setDocumentNumber(nextDocumentNumber(liabilitiesClaimsAssets, payment1));
payment2.setDocumentNumber(nextDocumentNumber(liabilitiesClaimsAssets, payment2));
newPaymentInstructions.add(payment1);
newPaymentInstructions.add(payment2);
} else {
// V / B
if (clearingCategory.contains(ClearingCategory.V)) {
log.debug("By liabilitiesClaimsAssets[{}] create for category V 1 record", liabilitiesClaimsAssets.getId());
PaymentInstruction payment1 = new PaymentInstruction();
payment1.setId(idGenerator.nextId());
payment1.setCreated(now);
payment1.setClearingDate(TimeUtil.toLocalDate(now));
payment1.setSenderId(SPVB_ID); // СПВБ
String sqlLCA = String.format("securityId = %s and companyId <> %s",
liabilitiesClaimsAssets.getSecurityId(), liabilitiesClaimsAssets.getCompanyId());
LiabilitiesClaimsAssets counterLiabilitiesClaimsAssets = liabilitiesClaimsAssetsImdg.getSingleObjectBySQL(sqlLCA);
if (counterLiabilitiesClaimsAssets == null) {
log.warn("Counter LiabilitiesClaimsAssets not found for : {}", sqlLCA);
} else {
payment1.setAddresseeId(counterLiabilitiesClaimsAssets.getCompanyId());
}
String symbol1 = selectSymbolValue(payment1.getAddresseeId(), CompanySymbol.BIC);
if (symbol1 == null) {
log.warn("CompanySymbols BIC not found for companyId={}", payment1.getAddresseeId());
} else {
payment1.setAdresseeBic(symbol1);
}
Company companyPRC = companyImdg.getSingleObjectByID(PRC_ID); // 2 "НКО АО ПРЦ"
String companyPRCName = null;
if (companyPRC == null) {
log.warn("Company.id=2 not found");
} else {
companyPRCName = companyPRC.getShortName();
}
payment1.setPayeeBankName(companyPRCName);
String symbolPRC = selectSymbolValue(PRC_ID, CompanySymbol.BIC); // 2 "НКО АО ПРЦ"
payment1.setPayeeBic(symbolPRC);
payment1.setAddresseeBankName(companyPRCName);
payment1.setPaymentDate(TimeUtil.localDateToInstant(liabilitiesClaimsAssets.getSettlementDate()));
payment1.setPaymentPurpose("Размещение депозита " + liabilitiesClaimsAssets.getContract());
payment1.setSettlementDate(liabilitiesClaimsAssets.getSettlementDate());
Long amount = liabilitiesClaimsAssets.getLiabilitiesQuantity() == null ? null : liabilitiesClaimsAssets.getLiabilitiesQuantity().abs().longValue();
payment1.setCreditLegAmount(amount);
payment1.setDebitLegAmount(amount);
{
String accSql = String.format("companyId = %s AND accountType=ANLT AND accountStatus=ACTV AND processingSign=ALWD",
payment1.getSenderId());
Account acc1 = accountImdg.getSingleObjectBySQL(accSql);
if (acc1 == null) {
log.warn("Account not found: {}", accSql);
} else {
payment1.setCreditLegAccountId(acc1.getId());
payment1.setCreditLegAccount(acc1.getAccount());
}
}
payment1.setCreditCsAccount(null);
{
Account acc1 = selectAccount(payment1.getAddresseeId(), AccountType.Corr, Status.Active, Allowed.ALLOWED);
if (acc1 != null) {
payment1.setDebitLegAccountId(acc1.getId());
payment1.setDebitLegAccount(acc1.getAccount());
}
}
payment1.setDebitCsAccount(null);
// payment1.setCreditLegDirection(InOutDirection.out.getKey());
// payment1.setDebitLegDirection(InOutDirection.in.getKey());
payment1.setCreditLegCurrencyCode(liabilitiesClaimsAssets.getCurrency());
payment1.setDebitLegCurrencyCode(liabilitiesClaimsAssets.getCurrency());
payment1.setTransactionStatus(TransactionStatus.stld.getKey());
payment1.setDocumentNumber(nextDocumentNumber(liabilitiesClaimsAssets, payment1));
newPaymentInstructions.add(payment1);
} else // Для V,B только 1 PaymentInstruction
if (clearingCategory.contains(ClearingCategory.B)) {
log.debug("By liabilitiesClaimsAssets[{}] create for category B 1 record", liabilitiesClaimsAssets.getId());
PaymentInstruction payment1 = new PaymentInstruction();
payment1.setId(idGenerator.nextId());
payment1.setCreated(now);
payment1.setClearingDate(TimeUtil.toLocalDate(now));
payment1.setSenderId(liabilitiesClaimsAssets.getCompanyId());
payment1.setAddresseeId(SPVB_ID); // СПВБ
String symbol1 = selectSymbolValue(payment1.getAddresseeId(), CompanySymbol.BIC);
if (symbol1 == null) {
log.warn("CompanySymbols BIC not found for companyId={}", payment1.getAddresseeId());
} else {
payment1.setAdresseeBic(symbol1);
}
Company companyPRC = companyImdg.getSingleObjectByID(PRC_ID); // 2 "НКО АО ПРЦ"
String companyPRCName = null;
if (companyPRC == null) {
log.warn("Company.id=2 not found");
} else {
companyPRCName = companyPRC.getShortName();
}
payment1.setPayeeBankName(companyPRCName);
String symbolPRC = selectSymbolValue(PRC_ID, CompanySymbol.BIC); // 2 "НКО АО ПРЦ"
payment1.setPayeeBic(symbolPRC);
payment1.setAddresseeBankName(companyPRCName);
payment1.setPaymentDate(TimeUtil.localDateToInstant(liabilitiesClaimsAssets.getSettlementDate()));
payment1.setPaymentPurpose("Размещение депозита " + liabilitiesClaimsAssets.getContract());
payment1.setSettlementDate(liabilitiesClaimsAssets.getSettlementDate());
Long amount = liabilitiesClaimsAssets.getLiabilitiesQuantity() == null ? null : liabilitiesClaimsAssets.getLiabilitiesQuantity().abs().longValue();
payment1.setCreditLegAmount(amount);
payment1.setDebitLegAmount(amount);
{
String accSql = String.format("companyId = %s AND accountType=ANLT AND accountStatus=ACTV AND processingSign=ALWD",
payment1.getSenderId());
Account acc1 = accountImdg.getSingleObjectBySQL(accSql);
if (acc1 == null) {
log.warn("Account not found: {}", accSql);
} else {
payment1.setCreditLegAccountId(acc1.getId());
payment1.setCreditLegAccount(acc1.getAccount());
}
}
payment1.setCreditCsAccount(null);
{
Account acc1 = selectAccount(payment1.getAddresseeId(), AccountType.Corr, Status.Active, Allowed.ALLOWED);
if (acc1 != null) {
payment1.setDebitLegAccountId(acc1.getId());
payment1.setDebitLegAccount(acc1.getAccount());
}
}
payment1.setDebitCsAccount(null);
// payment1.setCreditLegDirection(InOutDirection.out.getKey());
// payment1.setDebitLegDirection(InOutDirection.in.getKey());
payment1.setCreditLegCurrencyCode(liabilitiesClaimsAssets.getCurrency());
payment1.setDebitLegCurrencyCode(liabilitiesClaimsAssets.getCurrency());
payment1.setTransactionStatus(TransactionStatus.stld.getKey());
payment1.setDocumentNumber(nextDocumentNumber(liabilitiesClaimsAssets, payment1));
newPaymentInstructions.add(payment1);
}
}
return newPaymentInstructions;
}
protected String nextDocumentNumber(LiabilitiesClaimsAssets liabilitiesClaimsAssets, PaymentInstruction paymentInstruction) {
LocalDate nowD = LocalDate.now();
if (documentNumberResetAt == null || documentNumberResetAt.isBefore(nowD)) synchronized (this) {
long oldNum = documentNumberId.get(); // reset optimistic
if (oldNum > 0) {
while (!documentNumberId.compareAndSet(oldNum, 0)) {
oldNum = documentNumberId.get();
if (oldNum < 2) break;
}
}
documentNumberResetAt = nowD;
}
String paymentDate = DATE_FORMATTER_ddMMyy.format(TimeUtil.toLocalDate(paymentInstruction.getPaymentDate()));
String num = String.format("%s/%s/%s/%s",
liabilitiesClaimsAssets.getContract(), paymentDate,
paymentInstruction.getSenderId(), documentNumberId.incrementAndGet()
);
return num;
}
protected String selectSymbolValue(Long companyId, CompanySymbol symbol) {
CompanySymbols cSymbol = companySymbolsImdg.getSingleObjectByFieldValues(Map.of(
"companyId", companyId,
"companySymbol", symbol.getKey()));
if (cSymbol == null) {
return null;
} else {
return cSymbol.getCompanySymbolValue();
}
}
protected Account selectAccount(Long companyId, AccountType accountType, Status accountStatus, Allowed processingSign) {
Account account = accountImdg.getSingleObjectByFieldValues(Map.of(
"companyId", companyId,
"accountType", accountType.getKey(),
"accountStatus", accountStatus.getKey(),
"processingSign", processingSign.getKey()
));
return account;
}
/**
* 2.2.2.1. Выполнить Изменение liabilitiesClaimsAssets по paymentInstruction.
*/
public void updateFieldLiabilitiesClaimsAssetsByPayment(LiabilitiesClaimsAssets liabilitiesClaimsAssets,
List<PaymentInstruction> paymentInstructions,
Collection<ClearingCategory> clearingCategory) {
PaymentInstruction paymentIAny=paymentInstructions.iterator().next();
if (clearingCategory.contains(ClearingCategory.I) || clearingCategory.contains(ClearingCategory.V)) {
BigDecimal amount = paymentIAny.getCreditLegAmount() == null ? BigDecimal.ZERO : new BigDecimal((Long)paymentIAny.getCreditLegAmount());
// При обновлении (по paymentInstruction и категория I и счет CLRN):=liabilitiesClaimsAssets.liabilitiesQuantity - paymentInstruction.creditLeg_amount
// При обновлении (по paymentInstruction и категория V и счет INFO):=liabilitiesClaimsAssets.liabilitiesQuantity - paymentInstruction.creditLeg_amount
amount = liabilitiesClaimsAssets.getLiabilitiesQuantity().subtract(amount);
liabilitiesClaimsAssets.setLiabilitiesQuantity(amount);
}
if (clearingCategory.contains(ClearingCategory.B)) {
BigDecimal amount = paymentIAny.getCreditLegAmount() == null ? BigDecimal.ZERO : new BigDecimal((Long)paymentIAny.getCreditLegAmount());
// При обновлении (по paymentInstruction и категория B и счет CLRN): из paymentInstruction.creditLeg_amount
liabilitiesClaimsAssets.setLiabilitiesQuantity(amount);
}
if (clearingCategory.contains(ClearingCategory.I) || clearingCategory.contains(ClearingCategory.V)) {
BigDecimal amount = new BigDecimal((Long)paymentIAny.getDebitLegAmount());
liabilitiesClaimsAssets.setClaimsQuantity(amount);
}
if (clearingCategory.contains(ClearingCategory.B)) {
BigDecimal amount = new BigDecimal((Long)paymentIAny.getDebitLegAmount());
amount = liabilitiesClaimsAssets.getClaimsQuantity().subtract(amount);
liabilitiesClaimsAssets.setClaimsQuantity(amount);
}
// //на базе успешного/неуспешного изменения по paymentInstruction
// if (TransactionStatus.fail.equalsByKey(paymentIAny.getTransactionStatus())) {
// liabilitiesClaimsAssets.setClearingStatus(ClearingStatus.FAIL.getKey());
// } else if (TransactionStatus.ok.equalsByKey(paymentIAny.getTransactionStatus())) {
// liabilitiesClaimsAssets.setClearingStatus(ClearingStatus.OK.getKey());
// }
liabilitiesClaimsAssets.setUpdated(Instant.now());
}
/**
* 2.2.2.2. Выполнить Изменение liabilitiesClaimsMoney по paymentInstruction.
*/
public LiabilitiesClaimsMoney updateLiabilitiesClaimsMoneyByPaymentInstruction(LiabilitiesClaimsAssets liabilitiesClaimsAssets,
List<PaymentInstruction> paymentInstructions,
Collection<ClearingCategory> clearingCategory) {
Long lmId = liabilitiesClaimsAssets.getLiabilitiesClaimsMoneyId();
log.debug("Update LiabilitiesClaimsMoney[{}] by paymentInstructions {}", lmId, paymentInstructions.stream().map(SpcexObjectBase::getId).collect(Collectors.toList()));
LiabilitiesClaimsMoney liabilitiesClaimsMoney = liabilitiesClaimsMoneyImdg.getSingleObjectByID(lmId); // update data for update
if (liabilitiesClaimsMoney == null) {
throw new RuntimeException("IMDG not found LiabilitiesClaimsMoney.id=" + lmId);
}
if (!Objects.equals(liabilitiesClaimsMoney.getSettlementDate(), liabilitiesClaimsAssets.getSettlementDate())) {
log.warn("Date not match: LiabilitiesClaimsMoney[{}].settlementDate={} and LiabilitiesClaimsAssets[{}].settlementDate={}. It is first leg?",
liabilitiesClaimsMoney.getId(), liabilitiesClaimsMoney.getSettlementDate(), liabilitiesClaimsAssets.getId(), liabilitiesClaimsAssets.getSettlementDate()
);
}
updateFieldLiabilitiesClaimsMoneyByPayment(liabilitiesClaimsMoney, paymentInstructions, clearingCategory);
liabilitiesClaimsMoneyImdg.update(liabilitiesClaimsMoney);
return liabilitiesClaimsMoney;
}
/**
* Обновление liabilitiesClaimsMoney первой ноги по liabilitiesClaimsAssets
* @param liabilitiesClaimsMoney обновить этот объект (firstLeg)
* @param paymentInstructions
* @param clearingCategory
* @return были внесены изменения в liabilitiesClaimsMoney
*/
public boolean updateFieldLiabilitiesClaimsMoneyByPayment(LiabilitiesClaimsMoney liabilitiesClaimsMoney,
List<PaymentInstruction> paymentInstructions,
Collection<ClearingCategory> clearingCategory) {
boolean modified = false;
PaymentInstruction paymentIAny=paymentInstructions.iterator().next();
if ((clearingCategory.contains(ClearingCategory.I) || clearingCategory.contains(ClearingCategory.V))) {
BigDecimal amount = liabilitiesClaimsMoney.getLiabilitiesAmount().subtract(new BigDecimal((Long)paymentIAny.getCreditLegAmount()));
liabilitiesClaimsMoney.setLiabilitiesAmount(amount);
modified = true;
}
if ((clearingCategory.contains(ClearingCategory.B))) {
BigDecimal amount = liabilitiesClaimsMoney.getClaimsAmount().subtract(new BigDecimal((Long)paymentIAny.getDebitLegAmount()));
liabilitiesClaimsMoney.setClaimsAmount(amount);
modified = true;
}
if (modified) {
liabilitiesClaimsMoney.setUpdated(Instant.now());
}
return modified;
}
}

View file

@ -1,4 +1,4 @@
package ru.spcex.clearing.service.order;
package ru.spcex.clearing.service.builder;
import ru.clearing.classes.statics.data.payment.PaymentInstruction;
import ru.clearing.classes.statics.data.sdf.SDf03;

View file

@ -1,4 +1,4 @@
package ru.spcex.clearing.service.order;
package ru.spcex.clearing.service.builder;
import ru.clearing.classes.statics.data.payment.PaymentInstruction;
import ru.clearing.classes.statics.data.sdf.SDf11;

View file

@ -0,0 +1,111 @@
package ru.spcex.clearing.service.order;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import ru.clearing.classes.statics.data.execution.ExecutionDeposit;
import ru.spcex.platform.enumeration.ClearingCategory;
import ru.spcex.platform.enumeration.MoneyFlowSide;
import ru.spcex.platform.utils.collection.Pair;
import ru.spcex.platform.utils.enumeration.IEnumKey;
import java.util.*;
import java.util.function.BiConsumer;
import java.util.function.Consumer;
import java.util.function.Function;
import java.util.stream.Collectors;
@Component
public class ExecutionDepositSorter {
private Function<Long, ClearingCategory> clearingMemberCategoryProvider;
@Autowired
public void setClearingMemberCategoryProvider(Function<Long, ClearingCategory> clearingMemberCategoryProvider) {
this.clearingMemberCategoryProvider = clearingMemberCategoryProvider;
}
public LinkedHashMap<ClearingCategory, List<List<ExecutionDeposit>>> sortExecutionDeposit(Collection<ExecutionDeposit>execDeposits) {
//отсортировали все записи по категориям
// → внутри каждой подгруппы по категориям сортируем по компаниям
// → обрабатываем почередно все записи первой компании
// → когда по первой компании записи закончились, начинаем обработавать все записи второй компании и так далее
// → когда по первой категории закончились записи для всех компаниий,
// начинаем аналогично обработавать все записи второй категории и так далее).
//сортируем по категориям
Map<ClearingCategory, List<ExecutionDeposit>> categoryGrouping = execDeposits.stream().map(executionDeposit -> {
ClearingCategory category = clearingMemberCategoryProvider.apply(executionDeposit.getCompanyId());
return new Pair<>(category, executionDeposit);
}).collect(Collectors.groupingBy(Pair::getFirst, Collectors.mapping(Pair::getSecond, Collectors.toList())));
//внутри категорий сортируем по компаниям
LinkedHashMap<ClearingCategory, List<List<ExecutionDeposit>>> byCategoryByCompany = new LinkedHashMap<>();
Consumer<ClearingCategory> takeByClearingCategoryAndSortByCompanyId = clearingCategory -> {
List<ExecutionDeposit> wrappers = categoryGrouping.remove(clearingCategory);
if (wrappers != null) {
List<List<ExecutionDeposit>> batches = wrappers.stream()
.collect(Collectors.groupingBy(ExecutionDeposit::getCompanyId))
.entrySet()
.stream()
.sorted(Map.Entry.comparingByKey())
.map(Map.Entry::getValue)
.toList();
byCategoryByCompany.put(clearingCategory, batches);
}
};
//кладем в порядке I, V, B
takeByClearingCategoryAndSortByCompanyId.accept(ClearingCategory.I);
takeByClearingCategoryAndSortByCompanyId.accept(ClearingCategory.V);
takeByClearingCategoryAndSortByCompanyId.accept(ClearingCategory.B);
//для всех остальных порядок не важен, будут последними
categoryGrouping.keySet().forEach(takeByClearingCategoryAndSortByCompanyId);
//метод для дополнительной сортировки внутри групп по категории и компании
BiConsumer<ClearingCategory, Comparator<ExecutionDeposit>> additionalSorting =
(clearingCategory, comparator) -> {
List<List<ExecutionDeposit>> allExecDepositsForCategory = byCategoryByCompany.get(clearingCategory);
if (allExecDepositsForCategory != null) {
allExecDepositsForCategory.forEach(singleCompanyBatch -> singleCompanyBatch.sort(comparator));
}
};
//применяем дополнительные сортировки
//для V сначала Sell, потом Buy, в первую очередь должны обрабатываться новые сделки с наименьшими объемами
additionalSorting.accept(ClearingCategory.V,
firstSellThenBuy.thenComparing(ExecutionDeposit::getFirstLegAmount));
//для V сначала Buy, потом Sell, exchangeExecutionId - возрастающий
additionalSorting.accept(ClearingCategory.B,
firstBuyThenSell.thenComparing(ExecutionDeposit::getExchangeExecutionId));
return byCategoryByCompany;
}
public static Comparator<ExecutionDeposit> firstSellThenBuy = Comparator.comparing(wrapper -> {
MoneyFlowSide side = IEnumKey.getEnumByKey(MoneyFlowSide.class, wrapper.getSide());
if (side == null) return Integer.MAX_VALUE;
switch (side) {
case BUY -> {
return 1;
}
case SELL -> {
return 0;
}
default -> {
return Integer.MAX_VALUE;
}
}
});
public static Comparator<ExecutionDeposit> firstBuyThenSell = Comparator.comparing(wrapper -> {
MoneyFlowSide side = IEnumKey.getEnumByKey(MoneyFlowSide.class, wrapper.getSide());
if (side == null) return Integer.MAX_VALUE;
switch (side) {
case BUY -> {
return 0;
}
case SELL -> {
return 1;
}
default -> {
return Integer.MAX_VALUE;
}
}
});
}

View file

@ -1,7 +1,7 @@
package ru.spcex.clearing.service.order;
import ru.clearing.classes.statics.data.payment.PaymentInstruction;
import ru.spcex.platform.enumeration.ClearingMemberCategoryD;
import ru.spcex.platform.enumeration.ClearingCategory;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import java.util.Collection;
@ -13,13 +13,13 @@ public class PaymentBatchInfo {
private List<PaymentInstruction> initialOrder;
private List<PaymentInstruction> fromClearingToBank;
private List<PaymentInstruction> fromBankToClearing;
private ClearingMemberCategoryD categoryD;
private ClearingCategory categoryD;
private EnumMessage error;
public Stream<PaymentInstruction> getOrderedPaymentInstructions() {
Function<Collection<PaymentInstruction>, Stream<PaymentInstruction>> safeStream
= paymentInstructions -> paymentInstructions != null ? paymentInstructions.stream() : Stream.empty();
if (categoryD.equals(ClearingMemberCategoryD.B)) {
if (categoryD.equals(ClearingCategory.B)) {
return safeStream.apply(initialOrder);
} else {
return Stream.concat(safeStream.apply(fromClearingToBank), safeStream.apply(fromBankToClearing));
@ -43,11 +43,11 @@ public class PaymentBatchInfo {
this.fromBankToClearing = fromBankToClearing;
}
public ClearingMemberCategoryD getCategoryD() {
public ClearingCategory getCategoryD() {
return categoryD;
}
public void setCategoryD(ClearingMemberCategoryD categoryD) {
public void setCategoryD(ClearingCategory categoryD) {
this.categoryD = categoryD;
}

View file

@ -10,7 +10,7 @@ import ru.clearing.classes.statics.data.payment.PaymentInstruction;
import ru.spcex.clearing.error.ClearingError;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.platform.enumeration.AccountType;
import ru.spcex.platform.enumeration.ClearingMemberCategoryD;
import ru.spcex.platform.enumeration.ClearingCategory;
import ru.spcex.platform.enumeration.TransactionStatus;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
@ -46,15 +46,15 @@ public class PaymentInstructionSorter {
List<PaymentBatchInfo> result = new ArrayList<>(2);
log.info("processing PaymentInstruction's generationId={} senderId={} size={}", generationId, senderId, payments.size());
Collection<ClearingMemberCategory> cmcList = clrngMmbrImdg.getCollectionObjectsByFieldValues(Map.of("companyId", senderId));
Collection<ClearingMemberCategoryD> companyCategory = cmcList.stream()
Collection<ClearingCategory> companyCategory = cmcList.stream()
.filter(cat -> cat.getClearingMemberCategory() != null)
.map(cat -> IEnumKey.getEnumByKey(ClearingMemberCategoryD.class, cat.getClearingMemberCategory()))
.map(cat -> IEnumKey.getEnumByKey(ClearingCategory.class, cat.getClearingMemberCategory()))
.collect(Collectors.toSet());
for (ClearingMemberCategoryD categoryValue : companyCategory) {
for (ClearingCategory categoryValue : companyCategory) {
log.debug("For category {}", categoryValue);
PaymentBatchInfo batchInfo = new PaymentBatchInfo();
batchInfo.setCategoryD(categoryValue);
if (ClearingMemberCategoryD.I.equals(categoryValue)) {
if (ClearingCategory.I.equals(categoryValue)) {
EnumMessage error = null;
List<PaymentInstruction> aList = new ArrayList<>();
List<PaymentInstruction> bList = new ArrayList<>();
@ -118,7 +118,7 @@ public class PaymentInstructionSorter {
batchInfo.setFromBankToClearing(bList);
result.add(batchInfo);
}
} else if (ClearingMemberCategoryD.B.equals(categoryValue)) {
} else if (ClearingCategory.B.equals(categoryValue)) {
batchInfo.setInitialOrder(payments);
result.add(batchInfo);
} else {

View file

@ -0,0 +1,6 @@
package ru.spcex.clearing.service.validation;
public enum ClrngValidationStored {
relation, categoryPredefined
;
}

View file

@ -0,0 +1,89 @@
package ru.spcex.clearing.service.validation;
import ru.clearing.classes.statics.data.account.Account;
import ru.clearing.classes.statics.data.account.AccountBalance;
import ru.clearing.classes.statics.data.company.Company;
import ru.clearing.classes.statics.data.company.relation.Relation;
import ru.clearing.classes.statics.data.execution.ExecutionDeposit;
import ru.spcex.clearing.error.ClearingErrorInternal;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.platform.enumeration.AccountStatus;
import ru.spcex.platform.enumeration.ClearingCategory;
import ru.spcex.platform.enumeration.ServiceStatus;
import ru.spcex.platform.enumeration.WorkflowStatus;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.validation.ImdgValidationContext;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.validation.IValidationRule;
import java.util.Map;
import java.util.Optional;
public enum ExecutionDepositValidationRule implements IValidationRule<ImdgValidationContext<ExecutionDeposit>> {
ClearingIsAllowed() {
@Override
public Optional<EnumMessage> validate(ImdgValidationContext<ExecutionDeposit> context) {
ExecutionDeposit validatedObject = context.getValidatedObject();
Imdg<Relation> relationImdg = context.obtainMap(IMDGDistributedNames.Map_Relation, Relation.class);
Relation relation = relationImdg.getSingleObjectByFieldValues(Map.of("consumerId", validatedObject.getCompanyId()));
if (relation == null || !ServiceStatus.Active.equalsByKey(relation.getServiceStatus())) {
return of(ClearingErrorInternal.ClearingNotAllowed);
}
context.storeObject(ClrngValidationStored.relation, relation);
return empty();
}
},
AccountIsNotBlocked() {
@Override
public Optional<EnumMessage> validate(ImdgValidationContext<ExecutionDeposit> context) {
ExecutionDeposit validatedObject = context.getValidatedObject();
Imdg<Account> accountImdg = context.obtainMap(IMDGDistributedNames.Map_Account, Account.class);
Relation relation = context.getStoredObject(ClrngValidationStored.relation);
Account account = accountImdg.getSingleObjectByFieldValues(
Map.of("id", validatedObject.getAccountId(), "relationId", relation.getId()));
if (account == null || !AccountStatus.ACTIVE.equalsByKey(account.getAccountStatus())) {
return of(ClearingErrorInternal.AccountNotActive);
}
return empty();
}
},
CompanyIsNotBlocked() {
@Override
public Optional<EnumMessage> validate(ImdgValidationContext<ExecutionDeposit> context) {
ExecutionDeposit validatedObject = context.getValidatedObject();
Imdg<Company> companyImdg = context.obtainMap(IMDGDistributedNames.Map_Company, Company.class);
Company company = companyImdg.getSingleObjectByFieldValues(Map.of("id", validatedObject.getCompanyId()));
if (company == null || !WorkflowStatus.Active.equalsByKey(company.getWorkflowStatus())) {
return of(ClearingErrorInternal.CompanyNotActive);
}
return empty();
}
},
FinancialObligationSecurity() {
@Override
public Optional<EnumMessage> validate(ImdgValidationContext<ExecutionDeposit> context) {
ClearingCategory category = context.getStoredObject(ClrngValidationStored.categoryPredefined);
if (!category.equals(ClearingCategory.V)) {
return Optional.empty();
}
ExecutionDeposit validatedObject = context.getValidatedObject();
Imdg<AccountBalance> accountBalanceImdg = context.obtainMap(IMDGDistributedNames.Map_AccountBalance, AccountBalance.class);
AccountBalance accountBalance = accountBalanceImdg.getSingleObjectByFieldValues(
Map.of("accountId", validatedObject.getAccountId(),
"companyId", validatedObject.getCompanyId()));
if (accountBalance == null) {
return of(ClearingErrorInternal.FinancialObligationNotSatisfied);
}
//todo ABS
if (validatedObject.getFirstLegAmount().compareTo(accountBalance.getFreeBalanceAmount()) > 0) {
return of(ClearingErrorInternal.FinancialObligationNotSatisfied);
}
return empty();
}
};
@Override
public String ruleName() {
return "Sdf09ValidationRule." + name();
}
}

View file

@ -20,4 +20,3 @@ clearing-service.kafka-producer.linger-ms=1
clearing-service.kafka-producer.buffer-memory=33554432
clearing-service.scheduler.check-payment-instruction=*/5 * * * * *
clearing-service.scheduler.check-s-trade=0 1 0 1 * ?

View file

@ -34,4 +34,9 @@
<appender-ref ref="FILE"/>
<appender-ref ref="CONSOLE"/>
</logger>
<logger name="ru.spcex.clearing.service.Clearing" level="trace" additivity="false">
<appender-ref ref="FILE"/>
<appender-ref ref="CONSOLE"/>
</logger>
</configuration>

View file

@ -0,0 +1,29 @@
package ru.spcex.clearing.service.builder;
import org.junit.jupiter.api.Test;
import org.mockito.Mockito;
import ru.clearing.classes.statics.data.liabilities.LiabilitiesClaimsAssets;
import ru.clearing.classes.statics.data.payment.PaymentInstruction;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.utils.time.TimeUtil;
import java.time.LocalDate;
import static org.junit.jupiter.api.Assertions.assertEquals;
class PaymentInstructionCreatorTest {
@Test
void nextDocumentNumber() {
ImdgProvider nullImdg = Mockito.mock(ImdgProvider.class);
PaymentInstructionCreator instance = new PaymentInstructionCreator(nullImdg);
LiabilitiesClaimsAssets liabilitiesClaimsAssets = new LiabilitiesClaimsAssets();
liabilitiesClaimsAssets.setContract("contract");
PaymentInstruction paymentInstruction = new PaymentInstruction();
paymentInstruction.setPaymentDate(TimeUtil.localDateToInstant(LocalDate.of(2023,1,4)));
paymentInstruction.setSenderId(1234L);
assertEquals("contract/040123/1234/1", instance.nextDocumentNumber(liabilitiesClaimsAssets, paymentInstruction));
assertEquals("contract/040123/1234/2", instance.nextDocumentNumber(liabilitiesClaimsAssets, paymentInstruction));
assertEquals("contract/040123/1234/3", instance.nextDocumentNumber(liabilitiesClaimsAssets, paymentInstruction));
}
}

View file

@ -1,6 +1,8 @@
package ru.spcex.clearing.dbf.exporter.config;
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.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@ -14,6 +16,7 @@ import ru.spcex.platform.imdg.api.ImdgProvider;
//отдельный конфиг для sender чтобы сделать required false
@Configuration
public class KafkaSenderConfig {
Logger log = LoggerFactory.getLogger(getClass());
private final ImdgProvider imdgProvider;
private Producer<String, Object> kafkaProducer;
@ -23,14 +26,15 @@ public class KafkaSenderConfig {
this.imdgProvider = imdgProvider;
}
@Autowired
@Autowired(required = false)
public void setKafkaProducer(Producer<String, Object> kafkaProducer) {
this.kafkaProducer = kafkaProducer;
}
@Bean
public KafkaSender kafkaSender() {
if (kafkaProducer == null ||imdgProvider == null) {
if (kafkaProducer == null || imdgProvider == null) {
log.info("Can not create KafkaSender: kafkaProducer={}, imdgProvider={}", kafkaProducer, imdgProvider);
return null;
}
ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator();

View file

@ -27,6 +27,10 @@ public class ExecutionDepositMapStore extends TemplateMapStore<ExecutionDeposit>
return IMDGDistributedNames.Map_ExecutionDeposit;
}
public String[] getIndexingField() {
return new String[]{"exchangeExecutionId"};
}
@Override
public String[] getFields() {
return new String[]{"ID", "CREATED_AT", "UPDATED_AT",

View file

@ -109,6 +109,9 @@ public abstract class AbstractHazelcastLifecycleSupport implements InitializingB
}
} catch (InterruptedException | ExecutionException e) {
if (e instanceof InterruptedException) {
Thread.currentThread().interrupt();
}
throw new RuntimeException("MapStore multithreaded not complete.", e);
}

View file

@ -4,7 +4,7 @@ import com.thoughtworks.xstream.annotations.XStreamAlias;
import com.thoughtworks.xstream.annotations.XStreamAsAttribute;
import ru.spcex.clearing.reports.reports.ReportId;
import ru.spcex.clearing.reports.reports.ReportWithPeriod;
import ru.spcex.platform.enumeration.ClearingMemberCategoryD;
import ru.spcex.platform.enumeration.ClearingCategory;
import java.time.LocalDate;
import java.time.LocalDateTime;
@ -31,7 +31,7 @@ public class ReportPA_B extends ReportWithPeriod {
"Отчет по нетто-позиции участника клиринга в секции МКР (Уполномоченный банк)",
startDate,
endDate);
this.clearingMemberCategory = ClearingMemberCategoryD.B.getKey();
this.clearingMemberCategory = ClearingCategory.B.getKey();
}
@Override

View file

@ -4,7 +4,7 @@ import com.thoughtworks.xstream.annotations.XStreamAlias;
import com.thoughtworks.xstream.annotations.XStreamAsAttribute;
import ru.spcex.clearing.reports.reports.ReportId;
import ru.spcex.clearing.reports.reports.ReportWithPeriod;
import ru.spcex.platform.enumeration.ClearingMemberCategoryD;
import ru.spcex.platform.enumeration.ClearingCategory;
import java.time.LocalDate;
import java.time.LocalDateTime;
@ -33,8 +33,8 @@ public class ReportPA_IV extends ReportWithPeriod {
startDate,
endDate);
this.clearingMemberCategory = "%s,%s".formatted(
ClearingMemberCategoryD.I.getKey(),
ClearingMemberCategoryD.V.getKey()
ClearingCategory.I.getKey(),
ClearingCategory.V.getKey()
);
}

View file

@ -5,7 +5,7 @@ import com.thoughtworks.xstream.annotations.XStreamAsAttribute;
import ru.spcex.clearing.reports.reports.Destination;
import ru.spcex.clearing.reports.reports.ReportId;
import ru.spcex.clearing.reports.reports.ReportWithClearingStatus;
import ru.spcex.platform.enumeration.ClearingMemberCategoryD;
import ru.spcex.platform.enumeration.ClearingCategory;
import ru.spcex.platform.enumeration.ClearingStatus;
import java.time.LocalDate;
@ -37,8 +37,8 @@ public class ReportBR_0420315_P1 extends ReportWithClearingStatus {
ClearingStatus.OK.getKey());
this.destination = Destination.BR.getKey();
this.clearingMemberCategory = "%s,%s".formatted(
ClearingMemberCategoryD.I.getKey(),
ClearingMemberCategoryD.V.getKey()
ClearingCategory.I.getKey(),
ClearingCategory.V.getKey()
);
}

View file

@ -72,7 +72,7 @@ public abstract class AbstractReportCollector<T extends AbstractReport> {
*/
protected String toString(Double val) {
if (val == null) return "";
return String.format(Locale.US, "%.2f", val.doubleValue());
return String.format(Locale.US, "%.2f", val);
}
/**

View file

@ -8,7 +8,7 @@ import org.springframework.beans.factory.InitializingBean;
import org.springframework.context.annotation.Lazy;
import org.springframework.stereotype.Service;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.reports.ReportWithPeriodRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
import ru.spcex.clearing.reports.config.element.ReportsServiceSettings;
import ru.spcex.platform.enumeration.Task;
@ -40,20 +40,23 @@ public class QCommandExecutor extends QueueConsumer implements InitializingBean
@Override
public void afterPropertiesSet() {
log.info("Init queue listener {}", getClass().getSimpleName());
callback(ReportWithPeriodRequest.class)
callback(LauncherCommandRequest.class)
.setConsumer(this::newReportWithPeriod)
.forDestination(Task.createReport_GREP.topic(), callbacks::put);
init();
}
private void newReportWithPeriod(BaseRequest<ReportWithPeriodRequest> reportRequest) {
ReportWithPeriodRequest request = reportRequest.getRequestPayload();
private void newReportWithPeriod(BaseRequest<LauncherCommandRequest> reportRequest) {
LauncherCommandRequest request = reportRequest.getRequestPayload();
log.info("newReportWithPeriod request received: {}", request);
final String taskName = QUEUE_RUN_REPORT_COMMAND; //request.getTaskName();
if (StringUtils.isBlank(request.getReportId()))
throw new IllegalStateException("ReportID is empty");
LocalDate startDate = request.getStartDate();
LocalDate endDate = request.getEndDate();
// if (StringUtils.isBlank(request.getReportId()))
// throw new IllegalStateException("ReportID is empty");
// LocalDate startDate = request.getStartDate();
// LocalDate endDate = request.getEndDate();
LocalDate startDate = null;
LocalDate endDate = null;
if (startDate == null && endDate == null) {
endDate = LocalDate.now();
startDate = LocalDate.of(endDate.getYear(), endDate.getMonth(), 1);
@ -65,18 +68,18 @@ public class QCommandExecutor extends QueueConsumer implements InitializingBean
throw new IllegalStateException("StartDate > EndDate");
try {
if (QUEUE_RUN_REPORT_COMMAND.equals(request.getReportId())
|| QUEUE_RUN_DAILY_REPORT_COMMAND.equals(request.getReportId())) {
log.info("Generate all report by command {}", request.getReportId());
reportsServiceCommand.generateReportAll(request.getReportId(), startDate, endDate);
if (QUEUE_RUN_REPORT_COMMAND.equals(taskName)
|| QUEUE_RUN_DAILY_REPORT_COMMAND.equals(taskName)) {
log.info("Generate all report by command {}", taskName);
reportsServiceCommand.generateReportAll(taskName, startDate, endDate);
} else {
log.info("Generate single report by command {}", request.getReportId());
reportsServiceCommand.generateReportById(request.getReportId(), startDate, endDate);
log.info("Generate single report by command {}", taskName);
reportsServiceCommand.generateReportById(taskName, startDate, endDate);
}
log.debug("successfully processed");
} catch (IOException e) {
log.error("Error create report {}: {}", request.getReportId(), ExceptionUtils.getStackTrace(e));
log.error("Error create report {}: {}", taskName, ExceptionUtils.getStackTrace(e));
}
}

View file

@ -20,9 +20,9 @@ public abstract class ReportWithPeriodCollector<T extends AbstractReport> extend
LocalDate startDate,
LocalDate endDate) {
ImdgPredicateBuilder pb = liabilitiesClaimsAssetsImdg.predicateBuilder();
Instant dateTimeLo = toInstantStartDay(startDate);
Instant dateTimeHi = toInstantStartDay(endDate);
ImdgPredicate p = pb.and(pb.greatEqual("refundDate", dateTimeLo), pb.lessEqual("refundDate", dateTimeHi));
// Instant dateTimeLo = toInstantStartDay(startDate);
// Instant dateTimeHi = toInstantEndDay(endDate);
ImdgPredicate p = pb.and(pb.greatEqual("refundDate", startDate), pb.lessEqual("refundDate", endDate));
Collection<LiabilitiesClaimsAssets> liabilities = liabilitiesClaimsAssetsImdg.getCollectionObjectsByPredicate(p);
return liabilities;
}
@ -31,9 +31,9 @@ public abstract class ReportWithPeriodCollector<T extends AbstractReport> extend
LocalDate startDate, LocalDate endDate,
Long allowClearingStatusId) {
ImdgPredicateBuilder pb = liabilitiesClaimsAssetsImdg.predicateBuilder();
Instant dateTimeLo = toInstantStartDay(startDate);
Instant dateTimeHi = toInstantStartDay(endDate);
ImdgPredicate p = pb.and(pb.greatEqual("refundDate", dateTimeLo), pb.lessEqual("refundDate", dateTimeHi),
// Instant dateTimeLo = toInstantStartDay(startDate);
// Instant dateTimeHi = toInstantEndDay(endDate);
ImdgPredicate p = pb.and(pb.greatEqual("refundDate", startDate), pb.lessEqual("refundDate", endDate),
pb.equals("clearingStatus", allowClearingStatusId));
Collection<LiabilitiesClaimsAssets> liabilities = liabilitiesClaimsAssetsImdg.getCollectionObjectsByPredicate(p);
return liabilities;

View file

@ -12,6 +12,7 @@ import ru.spcex.clearing.reports.reports.ReportId;
import ru.spcex.clearing.reports.services.builders.XMLReportBuilder;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.utils.enumeration.IEnumKey;
import ru.spcex.platform.utils.log.ExceptionUtils;
import java.io.File;
import java.io.IOException;
@ -74,12 +75,20 @@ public class ReportsServiceCommand implements InitializingBean {
public void generateReportAll(@Nullable String reportGroup, LocalDate startDate, LocalDate endDate) throws IOException {
for (ReportWithPeriodCollector<?> reportDataCollector : allPeriodReports) {
if (reportGroup == null || QUEUE_RUN_REPORT_COMMAND.equals(reportGroup)) {
try {
makeSingleReport(reportDataCollector, startDate, endDate);
} catch (Exception e) {
log.error("Error generate report of type {}: {}", reportDataCollector, ExceptionUtils.getStackTrace(e));
}
}
}
for (SimpleReportCollector<?> reportDataCollector : allSimpleReports) {
if (reportGroup == null || QUEUE_RUN_DAILY_REPORT_COMMAND.equals(reportGroup)) {
try {
makeSingleReport(reportDataCollector);
} catch (Exception e) {
log.error("Error generate report of type {}: {}", reportDataCollector, ExceptionUtils.getStackTrace(e));
}
}
}
}

View file

@ -45,7 +45,7 @@ public class ReportBR_0420314_P2_Collector extends ReportWithPeriodCollector<Rep
Collection<PaymentInstruction> searchPaymentInstruction(LocalDate startDate, LocalDate endDate) {
ImdgPredicateBuilder pb = paymentInstructionImdg.predicateBuilder();
Instant dateTimeLo = toInstantStartDay(startDate);
Instant dateTimeHi = toInstantStartDay(endDate);
Instant dateTimeHi = toInstantEndDay(endDate);
ImdgPredicate p = pb.and(pb.greatEqual("created", dateTimeLo), pb.lessEqual("created", dateTimeHi));
Collection<PaymentInstruction> paymentInstructionsAll = paymentInstructionImdg.getCollectionObjectsByPredicate(p);
return paymentInstructionsAll;
@ -68,6 +68,10 @@ public class ReportBR_0420314_P2_Collector extends ReportWithPeriodCollector<Rep
final Long companyId = entry.getKey();
List<PaymentInstruction> paymentInstructionsByCompany = entry.getValue();
Company company = companyImdg.getSingleObjectByID(companyId);
if (company == null) {
log.warn("Company not found, id={}", companyId);
continue;
}
ProfileDocument profileDocumentCntr = profileDocumentImdg.getSingleObjectByFieldValues(Map.of(
"companyId", companyId,
"documentType", DocumentTypes.cntr.getKey()

View file

@ -35,7 +35,7 @@ public class ReportBR_0420317_Collector extends ReportWithPeriodCollector<Report
Collection<Relation> searchRelation(LocalDate startDate, LocalDate endDate, String service) {
ImdgPredicateBuilder pb = relationImdg.predicateBuilder();
Instant dateTimeLo = toInstantStartDay(startDate);
Instant dateTimeHi = toInstantStartDay(endDate);
Instant dateTimeHi = toInstantEndDay(endDate);
ImdgPredicate p = pb.and(pb.greatEqual("updated", dateTimeLo), pb.lessEqual("updated", dateTimeHi),
pb.equals("service", service));
Collection<Relation> relations = relationImdg.getCollectionObjectsByPredicate(p);

View file

@ -6,7 +6,7 @@ import ru.clearing.classes.statics.data.liabilities.LiabilitiesClaimsAssets;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.reports.reports.bt_12_2.ReportPA_B;
import ru.spcex.clearing.reports.services.ReportWithPeriodCollector;
import ru.spcex.platform.enumeration.ClearingMemberCategoryD;
import ru.spcex.platform.enumeration.ClearingCategory;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.imdg.api.predicate.ImdgPredicate;
@ -41,7 +41,7 @@ public class ReportPA_B_Collector extends ReportWithPeriodCollector<ReportPA_B>
Collection<ClearingMemberCategory> clearingMemberCategories =
clearingMemberCategoryImdg.getCollectionObjectsByFieldValues(Map.of(
"clearingMemberCategory", ClearingMemberCategoryD.B.getKey()
"clearingMemberCategory", ClearingCategory.B.getKey()
));
Set<Long> companyIds = clearingMemberCategories.stream()
.map(ClearingMemberCategory::getCompanyId)
@ -49,8 +49,8 @@ public class ReportPA_B_Collector extends ReportWithPeriodCollector<ReportPA_B>
ImdgPredicateBuilder builder = liabilitiesClaimsAssetsImdg.predicateBuilder();
ImdgPredicate refundDatePredicate = builder.and(
builder.greatEqual("refundDate", toInstantStartDay(startDate)),
builder.lessEqual("refundDate", toInstantStartDay(endDate))
builder.greatEqual("refundDate", (startDate)),
builder.lessEqual("refundDate", (endDate))
);
ImdgPredicate companyIdPredicates = builder.in("companyId", companyIds.toArray(new Comparable[0]));
ImdgPredicate finalPredicate = builder.and(refundDatePredicate, companyIdPredicates);

View file

@ -7,7 +7,7 @@ import ru.clearing.platform.dictionary.ClearingStatusDictionary;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.reports.reports.bt_12_2.ReportPA_IV;
import ru.spcex.clearing.reports.services.ReportWithPeriodCollector;
import ru.spcex.platform.enumeration.ClearingMemberCategoryD;
import ru.spcex.platform.enumeration.ClearingCategory;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.imdg.api.predicate.ImdgPredicate;
@ -45,7 +45,7 @@ public class ReportPA_IV_Collector extends ReportWithPeriodCollector<ReportPA_IV
validateDateRange(startDate, endDate);
String clearingMemberCategorySql = "clearingMemberCategory='%s' or clearingMemberCategory='%s'".formatted(
ClearingMemberCategoryD.I.getKey(), ClearingMemberCategoryD.V.getKey()
ClearingCategory.I.getKey(), ClearingCategory.V.getKey()
);
Collection<ClearingMemberCategory> clearingMemberCategories =
clearingMemberCategoryImdg.getCollectionObjectsBySQL(clearingMemberCategorySql);
@ -55,8 +55,8 @@ public class ReportPA_IV_Collector extends ReportWithPeriodCollector<ReportPA_IV
ImdgPredicateBuilder builder = liabilitiesClaimsAssetsImdg.predicateBuilder();
ImdgPredicate refundDatePredicate = builder.and(
builder.greatEqual("refundDate", toInstantStartDay(startDate)),
builder.lessEqual("refundDate", toInstantStartDay(endDate))
builder.greatEqual("refundDate", (startDate)),
builder.lessEqual("refundDate", (endDate))
);
ImdgPredicate companyIdPredicates = builder.in("companyId", companyIds.toArray(new Comparable[0]));
ImdgPredicate finalPredicate = builder.and(refundDatePredicate, companyIdPredicates);

View file

@ -14,7 +14,7 @@ class AbstractReportCollectorTest extends AbstractReportCollector {
@Test
void testToString() {
AbstractReportCollectorTest instance = new AbstractReportCollectorTest();
assertEquals(null, instance.toString((Long) null));
assertEquals("", instance.toString((Long) null));
assertEquals("1", instance.toString((Long) 1L));
assertEquals("1230", instance.toString((Long) 1230L));
@ -23,7 +23,7 @@ class AbstractReportCollectorTest extends AbstractReportCollector {
@Test
void testToString1() {
AbstractReportCollectorTest instance = new AbstractReportCollectorTest();
assertEquals(null, instance.toString((Instant) null));
assertEquals("", instance.toString((Instant) null));
Instant val = Instant.ofEpochMilli(1672058774338L);
assertEquals("26.12.2022T15:46:14", instance.toString(val)); // msk timezone
}
@ -38,7 +38,7 @@ class AbstractReportCollectorTest extends AbstractReportCollector {
@Test
void testToString3() {
AbstractReportCollectorTest instance = new AbstractReportCollectorTest();
assertEquals(null, instance.toString((BigDecimal)null));
assertEquals("", instance.toString((BigDecimal)null));
assertEquals("0.00", instance.toString(BigDecimal.ZERO));
assertEquals("1.10", instance.toString(BigDecimal.valueOf(1.10)));
assertEquals("100.12", instance.toString(BigDecimal.valueOf(100.12)));
@ -49,7 +49,7 @@ class AbstractReportCollectorTest extends AbstractReportCollector {
@Test
void testToString4() {
AbstractReportCollectorTest instance = new AbstractReportCollectorTest();
assertEquals(null, instance.toString((Double) null));
assertEquals("", instance.toString((Double) null));
assertEquals("0.00", instance.toString(0.0));
assertEquals("1.10", instance.toString(1.10));
assertEquals("100.12", instance.toString(100.12));

View file

@ -10,7 +10,6 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import ru.spcex.clearing.platform.messaging.domain.ActionType;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.reports.ReportWithPeriodRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
import ru.spcex.platform.enumeration.Task;
import ru.spcex.platform.imdg.api.ImdgId;
@ -33,16 +32,10 @@ public class LauncherSender {
protected BaseRequest<Object> makeCmdRequest(Task byTask, Long userId) {
Object toRequest;
if (Task.createReport_GREP.equals(byTask)) { // см. ru.spcex.clearing.backendapi.controller.request.cud.schedule.LauncherNew#toRequest
ReportWithPeriodRequest reportCommand = new ReportWithPeriodRequest();
reportCommand.setReportId(byTask.getKey());
toRequest = reportCommand;
} else {
LauncherCommandRequest taskRunnerCommandRequest = new LauncherCommandRequest();
taskRunnerCommandRequest.setTaskName(byTask.getKey());
taskRunnerCommandRequest.setUserId(userId);
toRequest = taskRunnerCommandRequest;
}
BaseRequest<Object> request = new BaseRequest<>();
request.setId(idGenerator.nextId());
request.setActionType(ActionType.NEW);

View file

@ -2,10 +2,12 @@ package ru.spcex.clearing.scheduler.service;
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.slf4j.Logger;
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.scheduler.Launcher;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
@ -20,14 +22,17 @@ import java.time.Instant;
import static ru.spcex.platform.utils.enumeration.IEnumKey.getEnumByKey;
@Service
public class LauncherService extends QueueConsumer implements InitializingBean {
private final Logger log = LoggerFactory.getLogger(getClass());
private final Imdg<Launcher> launcherMap;
private final Producer<String, Object> kafkaProducer;
@Autowired
public LauncherService(Consumer<String, Object> kafkaQueue, Producer<String, Object> kafkaProducer,
ImdgProvider imdgProvider) {
super(kafkaQueue, kafkaProducer);
this.kafkaProducer = kafkaProducer;
this.launcherMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Launcher, Launcher.class);
}
@ -53,6 +58,7 @@ public class LauncherService extends QueueConsumer implements InitializingBean {
launcher.setUpdated(created);
launcherMap.insert(launcher);
log.debug("successfully processed, new id {}", launcher.getId());
kafkaProducer.send(new ProducerRecord<>("launcher-" + launcher.getTask(), userRequest));
}
}

View file

@ -187,7 +187,6 @@ public class MoneyMarketSecurityService extends QueueConsumer implements Initial
}
MoneyMarketSecurity mms = validator.getStored(Stored.PresentById);
mms.setWorkflowStatus(ru.spcex.platform.enumeration.Status.Blocked.getKey());
Instant.now();
mms.setUpdated(Instant.now());
moneyMarketSecurityMap.update(mms);
Listing listing = listingImdg.getSingleObjectByFieldValues(Map.of("securityId", mms.getId()));

View file

@ -3,7 +3,7 @@ package ru.spcex.platform.enumeration;
import ru.spcex.platform.utils.enumeration.IEnumKey;
public enum AccountType implements IEnumKey {
Clrn("CLRN"), Bank("BANK"), Info("INFO");
Clrn("CLRN"), Bank("BANK"), Info("INFO"), Tran("TRAN"), Corr("CORR"), Anlt("ANLT");
private final String key;

View file

@ -3,7 +3,7 @@ package ru.spcex.platform.enumeration;
import ru.spcex.platform.utils.enumeration.IEnumKey;
public enum InOutDirection implements IEnumKey {
in("IN");
in("IN"), out("OUT");
private final String key;

View file

@ -2,12 +2,14 @@ package ru.spcex.platform.enumeration;
import ru.spcex.platform.utils.enumeration.IEnumKey;
public enum ClearingMemberCategoryD implements IEnumKey {
I("I"), V("V"), B("B");
public enum MoneyFlowSide implements IEnumKey {
BUY("BUY"),
SELL("SELL"),
;
private final String key;
ClearingMemberCategoryD(String key) {
MoneyFlowSide(String key) {
this.key = key;
}

View file

@ -0,0 +1,23 @@
package ru.spcex.platform.enumeration;
import ru.spcex.platform.utils.enumeration.IEnumKey;
public enum ServiceStatus implements IEnumKey {
Active("ACTV"),
Blocked("BLKD"),
Suspended("SSPD"),
Closed("CLOS"),
Reopened("ROPN"),
;
private final String key;
ServiceStatus(String key) {
this.key = key;
}
@Override
public String getKey() {
return key;
}
}

View file

@ -262,6 +262,7 @@ public abstract class HazelcastServiceBase
lock.wait();
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
log.error(ExceptionUtils.getStackTrace(e));
}
}

View file

@ -84,12 +84,14 @@ public final class HazelcastHelper {
log.info("Hazelcast: TextErrorService done");
done = true;
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new RuntimeException("Thread interrupted when waiting storage ready.", e);
} catch (Exception e) {
log.info(String.format("Waiting Hazelcast: %s -> %s", e.getClass().getSimpleName(), e.getMessage()));
try {
Thread.sleep(100);
} catch (InterruptedException ignored) {
Thread.currentThread().interrupt();
throw new RuntimeException("Thread interrupted.", e);
}
}

View file

@ -50,6 +50,8 @@ public interface Consts {
String EXPORT_PROCESS = "export-process";
String ACCOUNT_NEW = "account-new";
String BALANCE_ACCOUNT_NEW = "balance-account-new";
String BALANCE_ACCOUNT_UPDATE = "balance-account-update";
String CONTINUE_CLEARING = "continue-clearing";
String LAUNCHER_NEW = "launcher-new";
String NOTIFICATION_NEW = "notification-new";
@ -57,6 +59,7 @@ public interface Consts {
String REGISTRY_COVERED_DEAL_REGISTER_NEW = "registry-covered-deal-register-new";
String REGISTRY_ADMITTED_DEAL_REGISTER_NEW = "registry-admitted-deal-register-new";
String REGISTRY_DEAL_REGISTER_NEW = "registry-deal-register-new";
String REGISTRY_ORDER_REGISTER_NEW = "registry-order-register-new";
String REGISTRY_CONTRACT_REGISTER_NEW = "registry-contract-register-new";

View file

@ -0,0 +1,38 @@
package ru.spcex.clearing.platform.messaging.domain.cud.balance;
import com.fasterxml.jackson.annotation.JsonProperty;
import java.math.BigDecimal;
public class AccountBalanceClearingRequest {
@JsonProperty
private Long accountId;
@JsonProperty
private Long companyId;
@JsonProperty
private BigDecimal firstLegAmount;
public Long getAccountId() {
return accountId;
}
public void setAccountId(Long accountId) {
this.accountId = accountId;
}
public Long getCompanyId() {
return companyId;
}
public void setCompanyId(Long companyId) {
this.companyId = companyId;
}
public BigDecimal getFirstLegAmount() {
return firstLegAmount;
}
public void setFirstLegAmount(BigDecimal firstLegAmount) {
this.firstLegAmount = firstLegAmount;
}
}

View file

@ -0,0 +1,17 @@
package ru.spcex.clearing.platform.messaging.domain.cud.common;
import com.fasterxml.jackson.annotation.JsonProperty;
import ru.spcex.platform.classes.base.interfaces.WithId;
public class CommonIdRequest implements WithId {
@JsonProperty
public Long id;
public Long getId() {
return id;
}
public void setId(Long id) {
this.id = id;
}
}

View file

@ -0,0 +1,35 @@
package ru.spcex.clearing.platform.messaging.domain.cud.registry;
import com.fasterxml.jackson.annotation.JsonProperty;
import java.math.BigDecimal;
import java.time.Instant;
import java.time.LocalDate;
public class DealRegisterNewRequest {
@JsonProperty
public Long executionId;
@JsonProperty
public Long exchangeExecutionId;
public Long getExecutionId() {
return executionId;
}
public void setExecutionId(Long executionId) {
this.executionId = executionId;
}
public Long getExchangeExecutionId() {
return exchangeExecutionId;
}
public void setExchangeExecutionId(Long exchangeExecutionId) {
this.exchangeExecutionId = exchangeExecutionId;
}
}

View file

@ -6,9 +6,12 @@ public class LauncherCommandRequest {
@JsonProperty
private Long userId;
@JsonProperty
private String taskName;
@JsonProperty
private Long companyId;
@JsonProperty
private Long securityId;
public Long getUserId() {
@ -26,4 +29,20 @@ public class LauncherCommandRequest {
public void setTaskName(String taskName) {
this.taskName = taskName;
}
public Long getCompanyId() {
return companyId;
}
public void setCompanyId(Long companyId) {
this.companyId = companyId;
}
public Long getSecurityId() {
return securityId;
}
public void setSecurityId(Long securityId) {
this.securityId = securityId;
}
}

View file

@ -128,6 +128,9 @@ public class QueueConsumer implements AutoCloseable {
send = producer.send(new ProducerRecord<>(Consts.REQUEST_INFO_UPDATE, req));
send.get();
} catch (Exception e) {
if (e instanceof InterruptedException) {
Thread.currentThread().interrupt();
}
log.error(ExceptionUtils.getStackTrace(e));
}
});

View file

@ -45,6 +45,9 @@ public class KafkaSender {
try {
send.get();
} catch (InterruptedException | ExecutionException e) {
if (e instanceof InterruptedException) {
Thread.currentThread().interrupt();
}
log.error(ExceptionUtils.getStackTrace(e));
return null;
}

View file

@ -0,0 +1,27 @@
package ru.spcex.platform.utils.collection;
public class Pair<T1, T2> {
private T1 first;
private T2 second;
public Pair(T1 first, T2 second) {
this.first = first;
this.second = second;
}
public T1 getFirst() {
return first;
}
public void setFirst(T1 first) {
this.first = first;
}
public T2 getSecond() {
return second;
}
public void setSecond(T2 second) {
this.second = second;
}
}

View file

@ -61,4 +61,12 @@ public interface IEnumId extends Serializable {
default boolean equalsById(Long id) {
return id != null && getId().equals(id);
}
default String nameOrId() {
try {
return ((Enum<?>) this).name();
} catch (ClassCastException e) {
return getId().toString();
}
}
}

View file

@ -17,4 +17,9 @@ public class BigDecimalUtil {
return null;
}
}
public static BigDecimal safeBD(BigDecimal bd) {
return bd == null ? BigDecimal.ZERO : bd;
}
}

View file

@ -40,6 +40,7 @@
<folder_root_account-service>${folder_root_clearing}/clearing-parent/account-service</folder_root_account-service>
<folder_root_balance-service>${folder_root_clearing}/clearing-parent/balance-service</folder_root_balance-service>
<folder_root_company-service>${folder_root_clearing}/clearing-parent/company-service</folder_root_company-service>
<folder_root_clearing-service>${folder_root_clearing}/clearing-parent/clearing-service</folder_root_clearing-service>
<folder_root_reports-service>${folder_root_clearing}/clearing-parent/reports-service</folder_root_reports-service>
<folder_root_utility-service>${folder_root_clearing}/clearing-parent/utility-service</folder_root_utility-service>
<folder_root_scheduler-service>${folder_root_clearing}/clearing-parent/scheduler-service</folder_root_scheduler-service>

View file

@ -248,6 +248,25 @@
</fileSets>
</configuration>
</execution>
<execution>
<id>copy-clearing-service-bin</id>
<phase>prepare-package</phase>
<goals>
<goal>copy</goal>
</goals>
<configuration>
<fileSets>
<fileSet>
<sourceFile>${folder_root_clearing-service}/target/clearing-service.jar</sourceFile>
<destinationFile>${folder.clearing.distr.modules}/clearing-service/clearing-service.jar</destinationFile>
</fileSet>
<fileSet>
<sourceFile>${folder_root_clearing-service}/src/main/resources/application.properties</sourceFile>
<destinationFile>${folder.clearing.distr.modules}/clearing-service/application.properties</destinationFile>
</fileSet>
</fileSets>
</configuration>
</execution>
<execution>
<id>copy-scheduler-service-bin</id>
<phase>prepare-package</phase>