account-service http://jira.mfd.msk:8088/browse/CLS-271 поправил зацикливание и создание счёта по SDF-01
This commit is contained in:
parent
d8e519c838
commit
621069c933
8 changed files with 206 additions and 133 deletions
|
|
@ -100,9 +100,6 @@ public class AccountService extends QueueConsumer implements InitializingBean {
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public void afterPropertiesSet() {
|
public void afterPropertiesSet() {
|
||||||
callback(AccountSdf01Request.class)
|
|
||||||
.setFunction(this::accountNewSdf01)
|
|
||||||
.forDestination(Consts.ACCOUNT_NEW_SDF01, callbacks::put);
|
|
||||||
callback(CorrespondentAccountNewRequest.class)
|
callback(CorrespondentAccountNewRequest.class)
|
||||||
.setFunction(this::accountCorrespondentNew)
|
.setFunction(this::accountCorrespondentNew)
|
||||||
.forDestination(Consts.DESTINATION_CORRESPONDENT_ACCOUNT_NEW, callbacks::put);
|
.forDestination(Consts.DESTINATION_CORRESPONDENT_ACCOUNT_NEW, callbacks::put);
|
||||||
|
|
@ -234,50 +231,6 @@ public class AccountService extends QueueConsumer implements InitializingBean {
|
||||||
return null;
|
return null;
|
||||||
}
|
}
|
||||||
|
|
||||||
@Deprecated
|
|
||||||
public RequestInfoUpdate accountNewSdf01(BaseRequest<AccountSdf01Request> userRequest) {
|
|
||||||
log.debug("AccountSdf01Request received");
|
|
||||||
|
|
||||||
AccountSdf01Request req = userRequest.getRequestPayload();
|
|
||||||
List<AccountSdfToStatementRequestPart> accountToStatement = new ArrayList<>();
|
|
||||||
for (AccountSdfRequestPart accountReq : req.getAccounts()) {
|
|
||||||
Account account = new Account();
|
|
||||||
account.setAccount(accountReq.getAccount());
|
|
||||||
account.setCompanyId(accountReq.getCompanyId());
|
|
||||||
account.setAccountType(accountReq.getAccountType());
|
|
||||||
account.setStatus(WorkflowStatus.Active.getKey());
|
|
||||||
account.setCreated(Instant.now());
|
|
||||||
account.setUpdated(account.getCreated());
|
|
||||||
accountMap.insert(account);
|
|
||||||
AccountSdfToStatementRequestPart responsePart = responsePart(accountReq.getSdfId());
|
|
||||||
accountToStatement.add(responsePart);
|
|
||||||
}
|
|
||||||
sendStatementRequestBack(req.getGroupingSdf01Id(), accountToStatement);
|
|
||||||
log.debug("successfully processed, grouping id={}, processed number={}", req.getGroupingSdf01Id(), accountToStatement.size());
|
|
||||||
|
|
||||||
return null;
|
|
||||||
}
|
|
||||||
|
|
||||||
public void accountUpdateWithBrake(BaseRequest<AccountSdf01Request> userRequest) {
|
|
||||||
|
|
||||||
}
|
|
||||||
|
|
||||||
public AccountSdfToStatementRequestPart responsePart(Long sdf01Id) {
|
|
||||||
AccountSdfToStatementRequestPart responsePart = new AccountSdfToStatementRequestPart();
|
|
||||||
responsePart.setSdfId(sdf01Id);
|
|
||||||
responsePart.setErrorCode(null);
|
|
||||||
responsePart.setErrorText(null);
|
|
||||||
return responsePart;
|
|
||||||
}
|
|
||||||
|
|
||||||
@Deprecated
|
|
||||||
public void sendStatementRequestBack(Long groupingSdf01Id, List<AccountSdfToStatementRequestPart> results) {
|
|
||||||
StatementRequest request = new StatementRequest();
|
|
||||||
request.setGroupId(groupingSdf01Id);
|
|
||||||
request.setAccountCreationResults(results);
|
|
||||||
log.debug("Send message to kafka \"{}\": {}", Consts.STATEMENT_PROCESS, LogFormatter.toStringWrapper(request));
|
|
||||||
kafkaSender.sendRequestToQueue(Consts.STATEMENT_PROCESS, request);
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Заполняет поля relationId и companyId из соответствующей записи Relation
|
* Заполняет поля relationId и companyId из соответствующей записи Relation
|
||||||
|
|
|
||||||
|
|
@ -16,6 +16,10 @@ import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.Consts;
|
import ru.spcex.clearing.platform.messaging.domain.Consts;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.account.ClearingAccountNewRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.account.ClearingAccountNewRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.account.ClearingAccountUpdateRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.account.ClearingAccountUpdateRequest;
|
||||||
|
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.AccountSdfToStatementRequestPart;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.registry.TradingClearingRegistryNewRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.registry.TradingClearingRegistryNewRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.serialization.LogFormatter;
|
import ru.spcex.clearing.platform.messaging.serialization.LogFormatter;
|
||||||
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
|
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
|
||||||
|
|
@ -25,6 +29,7 @@ import ru.spcex.clearing.util.services.RequestHelper;
|
||||||
import ru.spcex.clearing.validation.common.ValidationHelper;
|
import ru.spcex.clearing.validation.common.ValidationHelper;
|
||||||
import ru.spcex.platform.enumeration.AccountType;
|
import ru.spcex.platform.enumeration.AccountType;
|
||||||
import ru.spcex.platform.enumeration.ServiceStatus;
|
import ru.spcex.platform.enumeration.ServiceStatus;
|
||||||
|
import ru.spcex.platform.enumeration.WorkflowStatus;
|
||||||
import ru.spcex.platform.imdg.api.Imdg;
|
import ru.spcex.platform.imdg.api.Imdg;
|
||||||
import ru.spcex.platform.imdg.api.ImdgProvider;
|
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||||
import ru.spcex.platform.imdg.api.ImdgTransaction;
|
import ru.spcex.platform.imdg.api.ImdgTransaction;
|
||||||
|
|
@ -34,6 +39,8 @@ import ru.spcex.platform.utils.enumeration.IMessageResolver;
|
||||||
import ru.spcex.platform.utils.validation.IValidator;
|
import ru.spcex.platform.utils.validation.IValidator;
|
||||||
|
|
||||||
import java.time.Instant;
|
import java.time.Instant;
|
||||||
|
import java.util.ArrayList;
|
||||||
|
import java.util.List;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
import java.util.function.Function;
|
import java.util.function.Function;
|
||||||
|
|
||||||
|
|
@ -61,9 +68,9 @@ public class ClearingAccountService extends QueueConsumer implements Initializin
|
||||||
IMessageResolver messageResolver,
|
IMessageResolver messageResolver,
|
||||||
RequestHelper requestHelper,
|
RequestHelper requestHelper,
|
||||||
@Qualifier("clearingAccountNewRequestValidator")
|
@Qualifier("clearingAccountNewRequestValidator")
|
||||||
Function<ClearingAccountNewRequest, IValidator> clearingAccountNewRequestValidator,
|
Function<ClearingAccountNewRequest, IValidator> clearingAccountNewRequestValidator,
|
||||||
@Qualifier("clearingAccountUpdateRequestValidator")
|
@Qualifier("clearingAccountUpdateRequestValidator")
|
||||||
Function<ClearingAccountUpdateRequest, IValidator> clearingAccountUpdateRequestValidator) {
|
Function<ClearingAccountUpdateRequest, IValidator> clearingAccountUpdateRequestValidator) {
|
||||||
super(kafkaQueue, kafkaResponseQueue);
|
super(kafkaQueue, kafkaResponseQueue);
|
||||||
this.kafkaSender = kafkaSender;
|
this.kafkaSender = kafkaSender;
|
||||||
this.accountService = accountService;
|
this.accountService = accountService;
|
||||||
|
|
@ -85,6 +92,10 @@ public class ClearingAccountService extends QueueConsumer implements Initializin
|
||||||
callback(ClearingAccountUpdateRequest.class)
|
callback(ClearingAccountUpdateRequest.class)
|
||||||
.setFunction(this::clearingAccountUpdate)
|
.setFunction(this::clearingAccountUpdate)
|
||||||
.forDestination(Consts.DESTINATION_CLEARING_ACCOUNT_UPDATE, callbacks::put);
|
.forDestination(Consts.DESTINATION_CLEARING_ACCOUNT_UPDATE, callbacks::put);
|
||||||
|
callback(AccountSdf01Request.class)
|
||||||
|
.setFunction(this::accountNewSdf01)
|
||||||
|
.forDestination(Consts.ACCOUNT_NEW_SDF01, callbacks::put);
|
||||||
|
|
||||||
init();
|
init();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -174,4 +185,104 @@ public class ClearingAccountService extends QueueConsumer implements Initializin
|
||||||
|
|
||||||
return null;
|
return null;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
public RequestInfoUpdate accountNewSdf01(BaseRequest<AccountSdf01Request> userRequest) {
|
||||||
|
log.debug("AccountSdf01Request received, id={}", userRequest.getId());
|
||||||
|
|
||||||
|
AccountSdf01Request req = userRequest.getRequestPayload();
|
||||||
|
List<AccountSdfToStatementRequestPart> accountToStatement = new ArrayList<>();
|
||||||
|
|
||||||
|
ImdgTransaction imdgTransaction = imdgProvider.newTransaction();
|
||||||
|
List<TradingClearingRegistryNewRequest> toTCRRequests = new ArrayList<>();
|
||||||
|
boolean txOk = false;
|
||||||
|
imdgTransaction.beginTransaction();
|
||||||
|
try {
|
||||||
|
accountsLoop:
|
||||||
|
for (AccountSdfRequestPart accountReq : req.getAccounts()) {
|
||||||
|
Instant now = Instant.now();
|
||||||
|
Account account = new Account();
|
||||||
|
account.setAccount(accountReq.getAccount());
|
||||||
|
account.setAccountType(AccountType.Clrn.getKey());
|
||||||
|
account.setStatus(ServiceStatus.Active.getKey());
|
||||||
|
account.setCompanyId(accountReq.getCompanyId());
|
||||||
|
account.setCreated(now);
|
||||||
|
account.setUpdated(now);
|
||||||
|
RequestInfoUpdate requestInfoUpdate = accountService.fillAccountFromRelation(account, userRequest.getId(), true);
|
||||||
|
if (requestInfoUpdate != null) {
|
||||||
|
log.warn("Error fill new account from relation. {}", /*account.getId(),*/ requestInfoUpdate.getMessage());
|
||||||
|
{
|
||||||
|
AccountSdfToStatementRequestPart responsePart = new AccountSdfToStatementRequestPart();
|
||||||
|
responsePart.setSdfId(accountReq.getSdfId());
|
||||||
|
responsePart.setErrorCode(AccountError.ClearingCategoryNotFound.getId()); // see accountService.fillAccountFromRelation
|
||||||
|
responsePart.setErrorText(requestInfoUpdate.getMessage());
|
||||||
|
accountToStatement.add(responsePart);
|
||||||
|
continue accountsLoop;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
Long clearingAccountId = -1L;
|
||||||
|
Long accountId = -1L;
|
||||||
|
ClearingAccount clearingAccount = null;
|
||||||
|
accountId = accountImdg.insert(account);
|
||||||
|
|
||||||
|
clearingAccount = new ClearingAccount();
|
||||||
|
clearingAccount.setCompanyId(accountReq.getCompanyId());
|
||||||
|
clearingAccount.setAccountId(accountId);
|
||||||
|
clearingAccount.setClearingAccountType(accountReq.getAccountType());
|
||||||
|
clearingAccountId = clearingAccountImdg.insert(clearingAccount);
|
||||||
|
log.trace("New account {}, clearingAccount {} was created.", accountId, clearingAccountId);
|
||||||
|
|
||||||
|
{ //todo как правильно заполнить реквест?
|
||||||
|
TradingClearingRegistryNewRequest request = new TradingClearingRegistryNewRequest();
|
||||||
|
if (AccountType.Depo.equalsByKey(clearingAccount.getClearingAccountType())) {
|
||||||
|
request.setDepoAccountId(accountId);
|
||||||
|
} else {
|
||||||
|
request.setMoneyAccountId(accountId);
|
||||||
|
}
|
||||||
|
request.setCompanyId(clearingAccount.getCompanyId());
|
||||||
|
toTCRRequests.add(request);
|
||||||
|
}
|
||||||
|
{
|
||||||
|
AccountSdfToStatementRequestPart responsePart = new AccountSdfToStatementRequestPart();
|
||||||
|
responsePart.setSdfId(accountReq.getSdfId());
|
||||||
|
responsePart.setErrorCode(null);
|
||||||
|
responsePart.setErrorText(null);
|
||||||
|
accountToStatement.add(responsePart);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
txOk = true;
|
||||||
|
} finally {
|
||||||
|
if (txOk) {
|
||||||
|
imdgTransaction.commitTransaction();
|
||||||
|
} else {
|
||||||
|
log.debug("failed insert, new clearing accounts. Request id={}", userRequest.getId());
|
||||||
|
imdgTransaction.rollbackTransaction();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if (txOk) {
|
||||||
|
log.debug("Sending {} messages of TradingClearingRegistryNewRequest", toTCRRequests.size());
|
||||||
|
for (TradingClearingRegistryNewRequest tcrReq : toTCRRequests) {
|
||||||
|
log.debug("Send message to kafka \"{}\": {}", Consts.DESTINATION_TRADING_CLEARING_REGISTRY_AUTO_NEW,
|
||||||
|
LogFormatter.toStringWrapper(tcrReq));
|
||||||
|
Long kafkaId = kafkaSender.sendRequestToQueue(Consts.DESTINATION_TRADING_CLEARING_REGISTRY_AUTO_NEW, tcrReq);
|
||||||
|
log.trace("successfully send request {} to kafka: new clearing account MoneyAccountId {}, DepoAccountId {}",
|
||||||
|
kafkaId, tcrReq.getMoneyAccountId(), tcrReq.getDepoAccountId());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
sendStatementRequestBack(req.getGroupingSdf01Id(), accountToStatement);
|
||||||
|
log.debug("successfully processed, grouping id={}, processed number={}", req.getGroupingSdf01Id(), accountToStatement.size());
|
||||||
|
|
||||||
|
return null;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void sendStatementRequestBack(Long groupingSdf01Id, List<AccountSdfToStatementRequestPart> results) {
|
||||||
|
StatementRequest request = new StatementRequest();
|
||||||
|
request.setGroupId(groupingSdf01Id);
|
||||||
|
request.setAccountCreationResults(results);
|
||||||
|
log.debug("Send message to kafka \"{}\": {}", Consts.STATEMENT_PROCESS, LogFormatter.toStringWrapper(request));
|
||||||
|
kafkaSender.sendRequestToQueue(Consts.STATEMENT_PROCESS, request);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -20,6 +20,7 @@ import ru.spcex.clearing.platform.messaging.domain.ActionType;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
|
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.Consts;
|
import ru.spcex.clearing.platform.messaging.domain.Consts;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.account.ClientCodeNewRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.account.ClientCodeNewRequest;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.CreateRegistryRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonIdRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonIdRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.registry.TradingClearingRegistryNewRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.registry.TradingClearingRegistryNewRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.registry.TradingClearingRegistryUpdateRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.registry.TradingClearingRegistryUpdateRequest;
|
||||||
|
|
@ -425,14 +426,11 @@ public class TradingClearingRegistryService extends QueueConsumer implements Ini
|
||||||
* clearing-service сообщение на открытие клиринговых регистров;
|
* clearing-service сообщение на открытие клиринговых регистров;
|
||||||
*/
|
*/
|
||||||
protected void sendNotificationToClearingSvc(TradingClearingRegistry tradingClearingRegistry) {
|
protected void sendNotificationToClearingSvc(TradingClearingRegistry tradingClearingRegistry) {
|
||||||
TradingClearingRegistryNewRequest request = new TradingClearingRegistryNewRequest();
|
CreateRegistryRequest request = new CreateRegistryRequest();
|
||||||
request.setTradingClearingRegistryType(tradingClearingRegistry.getTradingClearingRegistryType());
|
|
||||||
request.setCompanyId(tradingClearingRegistry.getCompanyId());
|
request.setCompanyId(tradingClearingRegistry.getCompanyId());
|
||||||
request.setMoneyAccountId(tradingClearingRegistry.getMoneyAccountId());
|
//todo Надо ли проверять существование регистров?
|
||||||
request.setDepoAccountId(tradingClearingRegistry.getDepoAccountId());
|
log.debug("Send message to kafka \"{}\": {}", Consts.REGISTRY_NEW, LogFormatter.toStringWrapper(request));
|
||||||
request.setStatus(tradingClearingRegistry.getStatus());
|
kafkaSender.sendRequestToQueue(Consts.REGISTRY_NEW, request);
|
||||||
log.debug("Send message to kafka \"{}\": {}", Consts.DESTINATION_TRADING_CLEARING_REGISTRY_AUTO_NEW, LogFormatter.toStringWrapper(request));
|
|
||||||
kafkaSender.sendRequestToQueue(Consts.DESTINATION_TRADING_CLEARING_REGISTRY_AUTO_NEW, request);
|
|
||||||
}
|
}
|
||||||
/**
|
/**
|
||||||
* todo report-service сообщение на формирование уведомления о создании нового ТКР
|
* todo report-service сообщение на формирование уведомления о создании нового ТКР
|
||||||
|
|
|
||||||
|
|
@ -70,7 +70,6 @@ class AccountServiceTest {
|
||||||
public static final MatcherFactory.Matcher<Account> ACCOUNT_MATCHER = usingIgnoringFieldsComparator();
|
public static final MatcherFactory.Matcher<Account> ACCOUNT_MATCHER = usingIgnoringFieldsComparator();
|
||||||
public static final MatcherFactory.Matcher<RequestInfo> REQUEST_INFO_MATCHER_MATCHER = usingIgnoringFieldsComparator("created");
|
public static final MatcherFactory.Matcher<RequestInfo> REQUEST_INFO_MATCHER_MATCHER = usingIgnoringFieldsComparator("created");
|
||||||
private static final int PARTITION = 0;
|
private static final int PARTITION = 0;
|
||||||
private static final String TOPIC_ACCOUNT_NEW = Consts.ACCOUNT_NEW_SDF01;
|
|
||||||
private static final String account = "123456789123";
|
private static final String account = "123456789123";
|
||||||
private static final Long companyId = 0L;
|
private static final Long companyId = 0L;
|
||||||
private static final Long relationId = 0L;
|
private static final Long relationId = 0L;
|
||||||
|
|
@ -247,76 +246,4 @@ class AccountServiceTest {
|
||||||
accountImdg.delete(resultBlock);
|
accountImdg.delete(resultBlock);
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
|
||||||
* {@link AccountService#accountNewSdf01(BaseRequest)}<br>
|
|
||||||
* Тест проверяет создание сущности {@link BaseRequest} в Hazelcast при передаче из Apache Kafka.<br>
|
|
||||||
* Входной запрос {@link AccountSdf01Request}:<br>
|
|
||||||
* {@link AccountSdfRequestPart#setSdfId} - текущий Id<br>
|
|
||||||
* {@link AccountSdfRequestPart#setAccount} - 123456789123<br>
|
|
||||||
* {@link AccountSdfRequestPart#setCompanyId} - текущий Id<br>
|
|
||||||
* {@link AccountSdf01Request#setGroupingSdf01Id} - текущий Id<br>
|
|
||||||
* {@link AccountSdf01Request#setAccounts} - Collections.singletonList(AccountSdfRequestPart)<br>
|
|
||||||
*/
|
|
||||||
@Test
|
|
||||||
void accountSdf01New() throws InterruptedException {
|
|
||||||
//ARRANGE
|
|
||||||
Long firstID = currentID.getAndIncrement();
|
|
||||||
Long secondID = currentID.getAndIncrement();
|
|
||||||
AccountSdfRequestPart accountSdfRequestPart = new AccountSdfRequestPart();
|
|
||||||
accountSdfRequestPart.setSdfId(firstID);
|
|
||||||
accountSdfRequestPart.setAccount(account);
|
|
||||||
accountSdfRequestPart.setCompanyId(firstID);
|
|
||||||
AccountSdf01Request accountSdf01Request = new AccountSdf01Request();
|
|
||||||
accountSdf01Request.setGroupingSdf01Id(firstID);
|
|
||||||
accountSdf01Request.setAccounts(Collections.singletonList(accountSdfRequestPart));
|
|
||||||
BaseRequest<AccountSdf01Request> baseNewRequest = new BaseRequest<>();
|
|
||||||
baseNewRequest.setRequestPayload(accountSdf01Request);
|
|
||||||
baseNewRequest.setId(firstID);
|
|
||||||
baseNewRequest.setActionType(ActionType.NEW);
|
|
||||||
String jsonBaseNewRequest;
|
|
||||||
ObjectMapper objectMapper = new ObjectMapper();
|
|
||||||
try {
|
|
||||||
jsonBaseNewRequest = objectMapper.writeValueAsString(baseNewRequest);
|
|
||||||
} catch (JsonProcessingException e) {
|
|
||||||
throw new RuntimeException(e);
|
|
||||||
}
|
|
||||||
|
|
||||||
AccountSdfToStatementRequestPart responsePart = new AccountSdfToStatementRequestPart();
|
|
||||||
responsePart.setSdfId(firstID);
|
|
||||||
responsePart.setErrorCode(null);
|
|
||||||
responsePart.setErrorText(null);
|
|
||||||
List<AccountSdfToStatementRequestPart> accountToStatement = Collections.singletonList(responsePart);
|
|
||||||
StatementRequest statementRequest = new StatementRequest();
|
|
||||||
statementRequest.setGroupId(firstID);
|
|
||||||
statementRequest.setAccountCreationResults(accountToStatement);
|
|
||||||
|
|
||||||
BaseRequest<Object> baseRequest = new BaseRequest<>();
|
|
||||||
baseRequest.setId(secondID);
|
|
||||||
baseRequest.setActionType(ActionType.SYSTEM);
|
|
||||||
baseRequest.setRequestPayload(statementRequest);
|
|
||||||
|
|
||||||
Account predictableAccount = new Account();
|
|
||||||
predictableAccount.setAccount(account);
|
|
||||||
predictableAccount.setCompanyId(firstID);
|
|
||||||
|
|
||||||
RequestInfo predictableRequestInfo = new RequestInfo();
|
|
||||||
predictableRequestInfo.setId(secondID);
|
|
||||||
predictableRequestInfo.setStatus(Status.Processing);
|
|
||||||
|
|
||||||
//KAFKA
|
|
||||||
addRecordToKafka((MockConsumer) accountService.getConsumer(), TOPIC_ACCOUNT_NEW, PARTITION, 0, jsonBaseNewRequest);
|
|
||||||
|
|
||||||
//waiting for kafka producer send message (finale event)
|
|
||||||
verify(producer, timeout(30_000L).times(2))
|
|
||||||
.send(producerRecord.capture());
|
|
||||||
//todo переписать валидацию ожидания на новые waitingSendAndCheckRecord / waitingWhenTryAddRecordAndCheckError
|
|
||||||
|
|
||||||
ImdgHazelcast<Account> accountImdg = (ImdgHazelcast<Account>) hazelcastServiceTest.getImdg(IMDGDistributedNames.Map_Account, Account.class);
|
|
||||||
|
|
||||||
//ASSERT
|
|
||||||
Account accountResult = accountImdg.getSingleObjectBySQL(String.format("account = %s", account));
|
|
||||||
predictableAccount.setId(accountResult.getId());
|
|
||||||
ACCOUNT_MATCHER.assertMatch(accountResult, predictableAccount);
|
|
||||||
accountImdg.delete(accountResult);
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
@ -1,5 +1,7 @@
|
||||||
package ru.spcex.clearing.account.service;
|
package ru.spcex.clearing.account.service;
|
||||||
|
|
||||||
|
import com.fasterxml.jackson.core.JsonProcessingException;
|
||||||
|
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||||
import org.apache.kafka.clients.consumer.MockConsumer;
|
import org.apache.kafka.clients.consumer.MockConsumer;
|
||||||
import org.apache.kafka.clients.producer.MockProducer;
|
import org.apache.kafka.clients.producer.MockProducer;
|
||||||
import org.apache.kafka.clients.producer.Producer;
|
import org.apache.kafka.clients.producer.Producer;
|
||||||
|
|
@ -26,23 +28,35 @@ import ru.spcex.clearing.account.config.validation.ClearingAccountValidationConf
|
||||||
import ru.spcex.clearing.account.config.validation.ValidationConfig;
|
import ru.spcex.clearing.account.config.validation.ValidationConfig;
|
||||||
import ru.spcex.clearing.account.utils.MatcherFactory;
|
import ru.spcex.clearing.account.utils.MatcherFactory;
|
||||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
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.Consts;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.account.ClearingAccountNewRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.account.ClearingAccountNewRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.account.ClearingAccountUpdateRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.account.ClearingAccountUpdateRequest;
|
||||||
|
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.AccountSdfToStatementRequestPart;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest;
|
||||||
|
import ru.spcex.clearing.platform.messaging.service.RequestInfo;
|
||||||
|
import ru.spcex.clearing.platform.messaging.service.Status;
|
||||||
import ru.spcex.clearing.test.TestObjectCreator;
|
import ru.spcex.clearing.test.TestObjectCreator;
|
||||||
import ru.spcex.clearing.test.config.ImdgTestConfig;
|
import ru.spcex.clearing.test.config.ImdgTestConfig;
|
||||||
import ru.spcex.clearing.test.config.KafkaTestConfig;
|
import ru.spcex.clearing.test.config.KafkaTestConfig;
|
||||||
import ru.spcex.platform.enumeration.*;
|
import ru.spcex.platform.enumeration.*;
|
||||||
import ru.spcex.platform.imdg.api.Imdg;
|
import ru.spcex.platform.imdg.api.Imdg;
|
||||||
|
import ru.spcex.platform.imdg.iml.hazelcast.adapter.ImdgHazelcast;
|
||||||
import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService;
|
import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService;
|
||||||
|
|
||||||
import javax.annotation.PostConstruct;
|
import javax.annotation.PostConstruct;
|
||||||
|
import java.util.Collections;
|
||||||
|
import java.util.List;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
|
|
||||||
import static org.mockito.Mockito.timeout;
|
import static org.mockito.Mockito.timeout;
|
||||||
import static org.mockito.Mockito.verify;
|
import static org.mockito.Mockito.verify;
|
||||||
import static ru.spcex.clearing.account.utils.MatcherFactory.usingIgnoringFieldsComparator;
|
import static ru.spcex.clearing.account.utils.MatcherFactory.usingIgnoringFieldsComparator;
|
||||||
import static ru.spcex.clearing.test.TestUtils.*;
|
import static ru.spcex.clearing.test.TestUtils.*;
|
||||||
|
import static ru.spcex.clearing.test.config.ImdgTestConfig.currentID;
|
||||||
|
|
||||||
@ExtendWith(SpringExtension.class)
|
@ExtendWith(SpringExtension.class)
|
||||||
@ContextConfiguration(classes = {
|
@ContextConfiguration(classes = {
|
||||||
|
|
@ -224,4 +238,77 @@ class ClearingAccountServiceTest {
|
||||||
ACCOUNT_MATCHER.assertMatch(accountResult, predictableAccount);
|
ACCOUNT_MATCHER.assertMatch(accountResult, predictableAccount);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// /**
|
||||||
|
// * {@link ClearingAccountService#accountNewSdf01(BaseRequest)}<br>
|
||||||
|
// * Тест проверяет создание сущности {@link BaseRequest} в Hazelcast при передаче из Apache Kafka.<br>
|
||||||
|
// * Входной запрос {@link AccountSdf01Request}:<br>
|
||||||
|
// * {@link AccountSdfRequestPart#setSdfId} - текущий Id<br>
|
||||||
|
// * {@link AccountSdfRequestPart#setAccount} - 123456789123<br>
|
||||||
|
// * {@link AccountSdfRequestPart#setCompanyId} - текущий Id<br>
|
||||||
|
// * {@link AccountSdf01Request#setGroupingSdf01Id} - текущий Id<br>
|
||||||
|
// * {@link AccountSdf01Request#setAccounts} - Collections.singletonList(AccountSdfRequestPart)<br>
|
||||||
|
// */
|
||||||
|
// @Test
|
||||||
|
// void accountSdf01New() throws InterruptedException {
|
||||||
|
// //ARRANGE
|
||||||
|
// Long firstID = currentID.getAndIncrement();
|
||||||
|
// Long secondID = currentID.getAndIncrement();
|
||||||
|
// AccountSdfRequestPart accountSdfRequestPart = new AccountSdfRequestPart();
|
||||||
|
// accountSdfRequestPart.setSdfId(firstID);
|
||||||
|
// accountSdfRequestPart.setAccount(account);
|
||||||
|
// accountSdfRequestPart.setCompanyId(firstID);
|
||||||
|
// AccountSdf01Request accountSdf01Request = new AccountSdf01Request();
|
||||||
|
// accountSdf01Request.setGroupingSdf01Id(firstID);
|
||||||
|
// accountSdf01Request.setAccounts(Collections.singletonList(accountSdfRequestPart));
|
||||||
|
// BaseRequest<AccountSdf01Request> baseNewRequest = new BaseRequest<>();
|
||||||
|
// baseNewRequest.setRequestPayload(accountSdf01Request);
|
||||||
|
// baseNewRequest.setId(firstID);
|
||||||
|
// baseNewRequest.setActionType(ActionType.NEW);
|
||||||
|
// String jsonBaseNewRequest;
|
||||||
|
// ObjectMapper objectMapper = new ObjectMapper();
|
||||||
|
// try {
|
||||||
|
// jsonBaseNewRequest = objectMapper.writeValueAsString(baseNewRequest);
|
||||||
|
// } catch (JsonProcessingException e) {
|
||||||
|
// throw new RuntimeException(e);
|
||||||
|
// }
|
||||||
|
//
|
||||||
|
// AccountSdfToStatementRequestPart responsePart = new AccountSdfToStatementRequestPart();
|
||||||
|
// responsePart.setSdfId(firstID);
|
||||||
|
// responsePart.setErrorCode(null);
|
||||||
|
// responsePart.setErrorText(null);
|
||||||
|
// List<AccountSdfToStatementRequestPart> accountToStatement = Collections.singletonList(responsePart);
|
||||||
|
// StatementRequest statementRequest = new StatementRequest();
|
||||||
|
// statementRequest.setGroupId(firstID);
|
||||||
|
// statementRequest.setAccountCreationResults(accountToStatement);
|
||||||
|
//
|
||||||
|
// BaseRequest<Object> baseRequest = new BaseRequest<>();
|
||||||
|
// baseRequest.setId(secondID);
|
||||||
|
// baseRequest.setActionType(ActionType.SYSTEM);
|
||||||
|
// baseRequest.setRequestPayload(statementRequest);
|
||||||
|
//
|
||||||
|
// Account predictableAccount = new Account();
|
||||||
|
// predictableAccount.setAccount(account);
|
||||||
|
// predictableAccount.setCompanyId(firstID);
|
||||||
|
//
|
||||||
|
// RequestInfo predictableRequestInfo = new RequestInfo();
|
||||||
|
// predictableRequestInfo.setId(secondID);
|
||||||
|
// predictableRequestInfo.setStatus(Status.Processing);
|
||||||
|
//
|
||||||
|
// //KAFKA
|
||||||
|
// final String TOPIC_ACCOUNT_NEW = Consts.ACCOUNT_NEW_SDF01;
|
||||||
|
// addRecordToKafka((MockConsumer) accountService.getConsumer(), TOPIC_ACCOUNT_NEW, PARTITION, 0, jsonBaseNewRequest);
|
||||||
|
//
|
||||||
|
// //waiting for kafka producer send message (finale event)
|
||||||
|
// verify(producer, timeout(30_000L).times(2))
|
||||||
|
// .send(producerRecord.capture());
|
||||||
|
// //todo переписать валидацию ожидания на новые waitingSendAndCheckRecord / waitingWhenTryAddRecordAndCheckError
|
||||||
|
//
|
||||||
|
// ImdgHazelcast<Account> accountImdg = (ImdgHazelcast<Account>) hazelcastServiceTest.getImdg(IMDGDistributedNames.Map_Account, Account.class);
|
||||||
|
//
|
||||||
|
// //ASSERT
|
||||||
|
// Account accountResult = accountImdg.getSingleObjectBySQL(String.format("account = %s", account));
|
||||||
|
// predictableAccount.setId(accountResult.getId());
|
||||||
|
// ACCOUNT_MATCHER.assertMatch(accountResult, predictableAccount);
|
||||||
|
// accountImdg.delete(accountResult);
|
||||||
|
// }
|
||||||
}
|
}
|
||||||
|
|
@ -85,7 +85,6 @@ public interface Consts {
|
||||||
String CREATE_REPORT_FOR_PERIOD = "create-report-for-period";
|
String CREATE_REPORT_FOR_PERIOD = "create-report-for-period";
|
||||||
String CREATE_REPORT_FOR_REGISTRY = "create-report-for-registry";
|
String CREATE_REPORT_FOR_REGISTRY = "create-report-for-registry";
|
||||||
|
|
||||||
@Deprecated
|
|
||||||
String ACCOUNT_NEW_SDF01 = "account-new-sdf01";
|
String ACCOUNT_NEW_SDF01 = "account-new-sdf01";
|
||||||
|
|
||||||
String DESTINATION_RELATION_NEW = "relation-new";
|
String DESTINATION_RELATION_NEW = "relation-new";
|
||||||
|
|
@ -100,7 +99,7 @@ public interface Consts {
|
||||||
String DESTINATION_CLIENT_CODE_UPDATE = "client-code-update";
|
String DESTINATION_CLIENT_CODE_UPDATE = "client-code-update";
|
||||||
String DESTINATION_CLIENT_CODE_DELETE = "client-code-delete";
|
String DESTINATION_CLIENT_CODE_DELETE = "client-code-delete";
|
||||||
|
|
||||||
String DESTINATION_TRADING_CLEARING_REGISTRY_NEW = "trading-clearing-registry-new"; // todo check buplicate REGISTRY_NEW?
|
String DESTINATION_TRADING_CLEARING_REGISTRY_NEW = "trading-clearing-registry-new"; // не путать с REGISTRY_NEW
|
||||||
String DESTINATION_TRADING_CLEARING_REGISTRY_AUTO_NEW = "trading-clearing-registry-auto-new";
|
String DESTINATION_TRADING_CLEARING_REGISTRY_AUTO_NEW = "trading-clearing-registry-auto-new";
|
||||||
String DESTINATION_TRADING_CLEARING_REGISTRY_UPDATE = "trading-clearing-registry-update";
|
String DESTINATION_TRADING_CLEARING_REGISTRY_UPDATE = "trading-clearing-registry-update";
|
||||||
String DESTINATION_TRADING_CLEARING_REGISTRY_BLOCK = "trading-clearing-registry-block";
|
String DESTINATION_TRADING_CLEARING_REGISTRY_BLOCK = "trading-clearing-registry-block";
|
||||||
|
|
|
||||||
|
|
@ -4,7 +4,6 @@ import com.fasterxml.jackson.annotation.JsonProperty;
|
||||||
|
|
||||||
import java.util.List;
|
import java.util.List;
|
||||||
|
|
||||||
@Deprecated
|
|
||||||
public class AccountSdf01Request {
|
public class AccountSdf01Request {
|
||||||
|
|
||||||
private Long groupingSdf01Id;
|
private Long groupingSdf01Id;
|
||||||
|
|
|
||||||
|
|
@ -2,7 +2,6 @@ package ru.spcex.clearing.platform.messaging.domain.cud.account.sdf01;
|
||||||
|
|
||||||
import com.fasterxml.jackson.annotation.JsonProperty;
|
import com.fasterxml.jackson.annotation.JsonProperty;
|
||||||
|
|
||||||
@Deprecated
|
|
||||||
public class AccountSdfRequestPart {
|
public class AccountSdfRequestPart {
|
||||||
@JsonProperty
|
@JsonProperty
|
||||||
private Long sdfId;
|
private Long sdfId;
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue