company-service/account-service http://jira.mfd.msk:8088/browse/CLS-272 рефакторинг: перенёс ClientCodeService.

This commit is contained in:
AKurakin 2023-05-16 18:05:58 +03:00
parent c1519a2d7a
commit 754df8edad
9 changed files with 9006 additions and 136 deletions

View file

@ -1,4 +1,4 @@
package ru.spcex.clearing.company.config.validation; package ru.spcex.clearing.account.config.validation;
import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Configuration;
@ -7,7 +7,7 @@ import ru.clearing.classes.statics.data.account.ClientCode;
import ru.clearing.classes.statics.data.company.Company; import ru.clearing.classes.statics.data.company.Company;
import ru.clearing.classes.statics.data.registry.TradingClearingRegistry; import ru.clearing.classes.statics.data.registry.TradingClearingRegistry;
import ru.clearing.platform.dictionary.WorkflowStatusDictionary; import ru.clearing.platform.dictionary.WorkflowStatusDictionary;
import ru.spcex.clearing.company.error.CompanyErrors; import ru.spcex.clearing.account.errors.AccountError;
import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.imdg.IMDGDistributedNames;
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.account.ClientCodeUpdateRequest; import ru.spcex.clearing.platform.messaging.domain.cud.account.ClientCodeUpdateRequest;
@ -40,33 +40,33 @@ public class ClientCodeValidationConfig {
addImdg.accept(IMDGDistributedNames.Map_TradingClearingRegistry); addImdg.accept(IMDGDistributedNames.Map_TradingClearingRegistry);
addImdg.accept(IMDGDistributedNames.Map_WorkflowStatusDictionary); addImdg.accept(IMDGDistributedNames.Map_WorkflowStatusDictionary);
return new ValidatorImpl<>(context, return new ValidatorImpl<>(context,
FieldRequiredRule.instance("companyId", ClientCodeNewRequest::getCompanyId, CompanyErrors.RequiredFieldEmpty), FieldRequiredRule.instance("companyId", ClientCodeNewRequest::getCompanyId, AccountError.RequiredFieldEmpty),
IdPresentRule.instance("companyId", IdPresentRule.instance("companyId",
ClientCodeNewRequest::getCompanyId, ClientCodeNewRequest::getCompanyId,
IMDGDistributedNames.Map_Company, IMDGDistributedNames.Map_Company,
Company.class, Company.class,
CompanyErrors.RequiredFieldEmpty, AccountError.RequiredFieldEmpty,
CompanyErrors.CompanyNotFound), AccountError.CompanyNotFound),
IdPresentRule.instance("moneyAccountId", IdPresentRule.instance("moneyAccountId",
ClientCodeNewRequest::getMoneyAccountId, ClientCodeNewRequest::getMoneyAccountId,
IMDGDistributedNames.Map_Account, IMDGDistributedNames.Map_Account,
Account.class, Account.class,
CompanyErrors.RequiredFieldEmpty, AccountError.RequiredFieldEmpty,
CompanyErrors.AccountNotFound, AccountError.AccountNotFound,
false, false,
acc -> Objects.equals(context.getValidatedObject().getCompanyId(), acc.getCompanyId()) acc -> Objects.equals(context.getValidatedObject().getCompanyId(), acc.getCompanyId())
? null : CompanyErrors.AccountNotFound ? null : AccountError.AccountNotFound
), ),
IdPresentRule.instance("depoAccountId", IdPresentRule.instance("depoAccountId",
ClientCodeNewRequest::getDepoAccountId, ClientCodeNewRequest::getDepoAccountId,
IMDGDistributedNames.Map_Account, IMDGDistributedNames.Map_Account,
Account.class, Account.class,
CompanyErrors.RequiredFieldEmpty, AccountError.RequiredFieldEmpty,
CompanyErrors.AccountNotFound, AccountError.AccountNotFound,
false, false,
acc -> Objects.equals(context.getValidatedObject().getCompanyId(), acc.getCompanyId()) acc -> Objects.equals(context.getValidatedObject().getCompanyId(), acc.getCompanyId())
? null : CompanyErrors.AccountNotFound ? null : AccountError.AccountNotFound
), ),
@ -74,8 +74,8 @@ public class ClientCodeValidationConfig {
ClientCodeNewRequest::getStatus, ClientCodeNewRequest::getStatus,
IMDGDistributedNames.Map_WorkflowStatusDictionary, IMDGDistributedNames.Map_WorkflowStatusDictionary,
WorkflowStatusDictionary.class, WorkflowStatusDictionary.class,
CompanyErrors.RequiredFieldEmpty, AccountError.RequiredFieldEmpty,
CompanyErrors.DictionaryNotFound, AccountError.DictionaryNotFound,
false) false)
); );
}; };
@ -97,45 +97,45 @@ public class ClientCodeValidationConfig {
ClientCodeUpdateRequest::getId, ClientCodeUpdateRequest::getId,
IMDGDistributedNames.Map_ClientCode, IMDGDistributedNames.Map_ClientCode,
ClientCode.class, ClientCode.class,
CompanyErrors.RequiredFieldEmpty, AccountError.RequiredFieldEmpty,
CompanyErrors.RecordNotFound AccountError.RecordNotFound
), ),
FieldRequiredRule.instance("companyId", ClientCodeUpdateRequest::getCompanyId, CompanyErrors.RequiredFieldEmpty), FieldRequiredRule.instance("companyId", ClientCodeUpdateRequest::getCompanyId, AccountError.RequiredFieldEmpty),
IdPresentRule.instance("companyId", IdPresentRule.instance("companyId",
ClientCodeUpdateRequest::getCompanyId, ClientCodeUpdateRequest::getCompanyId,
IMDGDistributedNames.Map_Company, IMDGDistributedNames.Map_Company,
Company.class, Company.class,
CompanyErrors.RequiredFieldEmpty, AccountError.RequiredFieldEmpty,
CompanyErrors.CompanyNotFound), AccountError.CompanyNotFound),
IdPresentRule.instance("moneyAccountId", IdPresentRule.instance("moneyAccountId",
ClientCodeUpdateRequest::getMoneyAccountId, ClientCodeUpdateRequest::getMoneyAccountId,
IMDGDistributedNames.Map_Account, IMDGDistributedNames.Map_Account,
Account.class, Account.class,
CompanyErrors.RequiredFieldEmpty, AccountError.RequiredFieldEmpty,
CompanyErrors.AccountNotFound, AccountError.AccountNotFound,
false, false,
acc -> Objects.equals(context.getValidatedObject().getCompanyId(), acc.getCompanyId()) acc -> Objects.equals(context.getValidatedObject().getCompanyId(), acc.getCompanyId())
? null : CompanyErrors.AccountNotFound ? null : AccountError.AccountNotFound
), ),
IdPresentRule.instance("depoAccountId", IdPresentRule.instance("depoAccountId",
ClientCodeUpdateRequest::getDepoAccountId, ClientCodeUpdateRequest::getDepoAccountId,
IMDGDistributedNames.Map_Account, IMDGDistributedNames.Map_Account,
Account.class, Account.class,
CompanyErrors.RequiredFieldEmpty, AccountError.RequiredFieldEmpty,
CompanyErrors.AccountNotFound, AccountError.AccountNotFound,
false, false,
acc -> Objects.equals(context.getValidatedObject().getCompanyId(), acc.getCompanyId()) acc -> Objects.equals(context.getValidatedObject().getCompanyId(), acc.getCompanyId())
? null : CompanyErrors.AccountNotFound ? null : AccountError.AccountNotFound
), ),
DictionaryPresentRule.instance("stauts", DictionaryPresentRule.instance("stauts",
ClientCodeUpdateRequest::getStatus, ClientCodeUpdateRequest::getStatus,
IMDGDistributedNames.Map_WorkflowStatusDictionary, IMDGDistributedNames.Map_WorkflowStatusDictionary,
WorkflowStatusDictionary.class, WorkflowStatusDictionary.class,
CompanyErrors.RequiredFieldEmpty, AccountError.RequiredFieldEmpty,
CompanyErrors.DictionaryNotFound, AccountError.DictionaryNotFound,
false) false)
); );
}; };
@ -154,8 +154,8 @@ public class ClientCodeValidationConfig {
CommonDeleteRequest::getId, CommonDeleteRequest::getId,
IMDGDistributedNames.Map_ClientCode, IMDGDistributedNames.Map_ClientCode,
ClientCode.class, ClientCode.class,
CompanyErrors.RequiredFieldEmpty, AccountError.RequiredFieldEmpty,
CompanyErrors.RecordNotFound) AccountError.RecordNotFound)
); );
}; };
} }

View file

@ -7,10 +7,7 @@ import ru.clearing.classes.statics.data.account.*;
import ru.clearing.classes.statics.data.company.ClearingMemberCategory; import ru.clearing.classes.statics.data.company.ClearingMemberCategory;
import ru.clearing.classes.statics.data.company.Company; import ru.clearing.classes.statics.data.company.Company;
import ru.clearing.classes.statics.data.registry.TradingClearingRegistry; import ru.clearing.classes.statics.data.registry.TradingClearingRegistry;
import ru.clearing.platform.dictionary.AccountTypeDictionary; import ru.clearing.platform.dictionary.*;
import ru.clearing.platform.dictionary.ClearingAccountTypeDictionary;
import ru.clearing.platform.dictionary.CurrencyCodeDictionary;
import ru.clearing.platform.dictionary.ServiceStatusDictionary;
import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.validation.common.ValidationHelper; import ru.spcex.clearing.validation.common.ValidationHelper;
import ru.spcex.platform.classes.base.SpcexObjectBase; import ru.spcex.platform.classes.base.SpcexObjectBase;
@ -43,6 +40,10 @@ public class ValidationConfig {
addImdg.accept(IMDGDistributedNames.Map_ServiceStatusDictionary, ServiceStatusDictionary.class); addImdg.accept(IMDGDistributedNames.Map_ServiceStatusDictionary, ServiceStatusDictionary.class);
addImdg.accept(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class); addImdg.accept(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class);
//for ClientCodeValidationConfig
addImdg.accept(IMDGDistributedNames.Map_WorkflowStatusDictionary, WorkflowStatusDictionary.class);
addImdg.accept(IMDGDistributedNames.Map_ClientCode, ClientCode.class);
return imdg; return imdg;
} }

View file

@ -6,7 +6,9 @@ public enum AccountError implements IErrorEnumId {
GeneralError(5000L), GeneralError(5000L),
UserVerifyDenial(5001L), UserVerifyDenial(5001L),
RequiredFieldEmpty(5002L), RequiredFieldEmpty(5002L),
DictionaryNotFound(5003L),
WrongFieldValue(5004L), WrongFieldValue(5004L),
RecordNotFound(5006L),
AccountAlreadyExist(5010L), AccountAlreadyExist(5010L),
AccountNotFound(5011L), AccountNotFound(5011L),
AccountNotActive(5012L), AccountNotActive(5012L),

View file

@ -1,4 +1,4 @@
package ru.spcex.clearing.company.service; package ru.spcex.clearing.account.service;
import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.producer.Producer; import org.apache.kafka.clients.producer.Producer;
@ -14,7 +14,7 @@ import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.account.ClientCode; import ru.clearing.classes.statics.data.account.ClientCode;
import ru.clearing.classes.statics.data.company.Company; import ru.clearing.classes.statics.data.company.Company;
import ru.clearing.classes.statics.data.registry.TradingClearingRegistry; import ru.clearing.classes.statics.data.registry.TradingClearingRegistry;
import ru.spcex.clearing.company.error.CompanyErrors; import ru.spcex.clearing.account.errors.AccountError;
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.ActionType;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest; import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
@ -26,8 +26,10 @@ import ru.spcex.clearing.platform.messaging.domain.cud.registry.TradingClearingR
import ru.spcex.clearing.platform.messaging.domain.cud.registry.TradingClearingRegistryUpdateRequest; import ru.spcex.clearing.platform.messaging.domain.cud.registry.TradingClearingRegistryUpdateRequest;
import ru.spcex.clearing.platform.messaging.service.QueueConsumer; import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate; import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate;
import ru.spcex.clearing.platform.messaging.service.Status;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.clearing.util.security.UserRoleVerification; import ru.spcex.clearing.util.security.UserRoleVerification;
import ru.spcex.clearing.util.services.RequestHelper;
import ru.spcex.clearing.util.services.exchangers.BiDirectionQueueExchanger; import ru.spcex.clearing.util.services.exchangers.BiDirectionQueueExchanger;
import ru.spcex.clearing.validation.common.ValidationHelper; import ru.spcex.clearing.validation.common.ValidationHelper;
import ru.spcex.platform.enumeration.TradingClearingRegistryType; import ru.spcex.platform.enumeration.TradingClearingRegistryType;
@ -51,12 +53,11 @@ import java.util.concurrent.Future;
import java.util.function.Function; import java.util.function.Function;
@Service @Service
public class ClientCodeService extends QueueConsumer implements InitializingBean, DisposableBean { public class ClientCodeService extends QueueConsumer implements InitializingBean {
private final Logger log = LoggerFactory.getLogger(getClass()); private final Logger log = LoggerFactory.getLogger(getClass());
private final Producer<String, Object> kafkaProducer; private final Producer<String, Object> kafkaProducer;
private final ImdgId idGenerator; private final ImdgId idGenerator;
private final Imdg<ClientCode> clientCodeMap; private final Imdg<ClientCode> clientCodeMap;
private final Imdg<Company> companyMap;
private final Imdg<TradingClearingRegistry> tradingClearingRegistryMap; private final Imdg<TradingClearingRegistry> tradingClearingRegistryMap;
@ -66,37 +67,35 @@ public class ClientCodeService extends QueueConsumer implements InitializingBean
private final ValidationHelper validationHelper; private final ValidationHelper validationHelper;
private final UserRoleVerification userRoleVerification; private final UserRoleVerification userRoleVerification;
private final IMessageResolver messageResolver; private final IMessageResolver messageResolver;
private final RequestHelper requestHelper;
BiDirectionQueueExchanger<BaseRequest<TradingClearingRegistryNewRequest>> accountServiceExchanger; protected TradingClearingRegistryService tradingClearingRegistryService;
@Autowired @Autowired
public ClientCodeService(Consumer<String, Object> kafkaQueue1, Consumer<String, Object> kafkaQueue2, public ClientCodeService(Consumer<String, Object> kafkaQueue,
Producer<String, Object> kafkaProducer, Producer<String, Object> kafkaProducer,
ImdgProvider imdgProvider, ImdgProvider imdgProvider,
ValidationHelper validationHelper, ValidationHelper validationHelper,
IMessageResolver messageResolver, IMessageResolver messageResolver,
RequestHelper requestHelper,
TradingClearingRegistryService tradingClearingRegistryService,
UserRoleVerification userRoleVerification, UserRoleVerification userRoleVerification,
@Qualifier("clientCodeNewRequestValidator") Function<ClientCodeNewRequest, IValidator> clientCodeNewRequestValidator, @Qualifier("clientCodeNewRequestValidator") Function<ClientCodeNewRequest, IValidator> clientCodeNewRequestValidator,
@Qualifier("clientCodeUpdateRequestValidator") Function<ClientCodeUpdateRequest, IValidator> clientCodeUpdateRequestValidator, @Qualifier("clientCodeUpdateRequestValidator") Function<ClientCodeUpdateRequest, IValidator> clientCodeUpdateRequestValidator,
@Qualifier("clientCodeDeleteRequestValidator") Function<CommonDeleteRequest, IValidator> clientCodeDeleteRequestValidator) { @Qualifier("clientCodeDeleteRequestValidator") Function<CommonDeleteRequest, IValidator> clientCodeDeleteRequestValidator) {
super(kafkaQueue1, kafkaProducer); super(kafkaQueue, kafkaProducer);
this.kafkaProducer = kafkaProducer; this.kafkaProducer = kafkaProducer;
this.idGenerator = imdgProvider.getImdgIdGenerator(); this.idGenerator = imdgProvider.getImdgIdGenerator();
this.clientCodeMap = imdgProvider.getImdg(IMDGDistributedNames.Map_ClientCode, ClientCode.class); this.clientCodeMap = imdgProvider.getImdg(IMDGDistributedNames.Map_ClientCode, ClientCode.class);
this.companyMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Company, Company.class);
this.tradingClearingRegistryMap = imdgProvider.getImdg(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class); this.tradingClearingRegistryMap = imdgProvider.getImdg(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class);
this.clientCodeNewRequestValidator = clientCodeNewRequestValidator; this.clientCodeNewRequestValidator = clientCodeNewRequestValidator;
this.clientCodeUpdateRequestValidator = clientCodeUpdateRequestValidator; this.clientCodeUpdateRequestValidator = clientCodeUpdateRequestValidator;
this.clientCodeDeleteRequestValidator = clientCodeDeleteRequestValidator; this.clientCodeDeleteRequestValidator = clientCodeDeleteRequestValidator;
this.validationHelper = validationHelper; this.validationHelper = validationHelper;
this.requestHelper = requestHelper;
this.messageResolver = messageResolver; this.messageResolver = messageResolver;
this.userRoleVerification = userRoleVerification; this.userRoleVerification = userRoleVerification;
this.tradingClearingRegistryService = tradingClearingRegistryService;
this.accountServiceExchanger = new BiDirectionQueueExchanger<>(kafkaQueue2, kafkaProducer,
Consts.DESTINATION_TRADING_CLEARING_REGISTRY_NEW,
Consts.DESTINATION_TRADING_CLEARING_REGISTRY_REPLY, true, // REQUEST_INFO_UPDATE - стандартная очередь, для результатов всех реквестов. DESTINATION_TRADING_CLEARING_REGISTRY_REPLY,
60000
);
} }
@Override @Override
@ -118,10 +117,6 @@ public class ClientCodeService extends QueueConsumer implements InitializingBean
} }
@Override
public void destroy() {
accountServiceExchanger.close();
}
protected RequestInfoUpdate clientCodeNew(BaseRequest<ClientCodeNewRequest> userRequest) { protected RequestInfoUpdate clientCodeNew(BaseRequest<ClientCodeNewRequest> userRequest) {
log.debug("ClientCodeNewRequest received {}", userRequest.getId()); log.debug("ClientCodeNewRequest received {}", userRequest.getId());
@ -143,7 +138,7 @@ public class ClientCodeService extends QueueConsumer implements InitializingBean
log.error("Can not wait creation of TCR. request id={};CompanyId={}, MoneyAccountId={}, DepoAccountId={}; {}", log.error("Can not wait creation of TCR. request id={};CompanyId={}, MoneyAccountId={}, DepoAccountId={}; {}",
userRequest.getId(), req.getCompanyId(), req.getMoneyAccountId(), req.getDepoAccountId(), userRequest.getId(), req.getCompanyId(), req.getMoneyAccountId(), req.getDepoAccountId(),
e.toString()); e.toString());
return makeErrorResponse(userRequest, CompanyErrors.GeneralError, "Can not create TCR at this moment"); return makeErrorResponse(userRequest, AccountError.GeneralError, "Can not create TCR: " + e.getMessage());
} }
} }
@ -176,7 +171,7 @@ public class ClientCodeService extends QueueConsumer implements InitializingBean
log.error("Can not wait creation of TCR. request id={};CompanyId={}, MoneyAccountId={}, DepoAccountId={}; {}", log.error("Can not wait creation of TCR. request id={};CompanyId={}, MoneyAccountId={}, DepoAccountId={}; {}",
userRequest.getId(), req.getCompanyId(), req.getMoneyAccountId(), req.getDepoAccountId(), userRequest.getId(), req.getCompanyId(), req.getMoneyAccountId(), req.getDepoAccountId(),
e.toString()); e.toString());
return makeErrorResponse(userRequest, CompanyErrors.GeneralError, "Can not create TCR at this moment"); return makeErrorResponse(userRequest, AccountError.GeneralError, "Can not create TCR at this moment" + e.getMessage());
} }
} }
@ -201,7 +196,7 @@ public class ClientCodeService extends QueueConsumer implements InitializingBean
log.debug("ClientCodeUpdateRequest received, id={}", req.getId()); log.debug("ClientCodeUpdateRequest received, id={}", req.getId());
ClientCode clientCode = clientCodeMap.getSingleObjectByID(req.getId()); ClientCode clientCode = clientCodeMap.getSingleObjectByID(req.getId());
if (clientCode == null) { if (clientCode == null) {
return makeErrorResponse(userRequest, CompanyErrors.RecordNotFound, req.getId()); return makeErrorResponse(userRequest, AccountError.RecordNotFound, req.getId());
} }
boolean doCreateTCR = checkNeedCreateTCR(req.getCompanyId(), req.getMoneyAccountId(), req.getDepoAccountId()); boolean doCreateTCR = checkNeedCreateTCR(req.getCompanyId(), req.getMoneyAccountId(), req.getDepoAccountId());
@ -213,7 +208,7 @@ public class ClientCodeService extends QueueConsumer implements InitializingBean
log.error("Can not wait creation of TCR. request id={};CompanyId={}, MoneyAccountId={}, DepoAccountId={}; {}", log.error("Can not wait creation of TCR. request id={};CompanyId={}, MoneyAccountId={}, DepoAccountId={}; {}",
userRequest.getId(), req.getCompanyId(), req.getMoneyAccountId(), req.getDepoAccountId(), userRequest.getId(), req.getCompanyId(), req.getMoneyAccountId(), req.getDepoAccountId(),
e.toString()); e.toString());
return makeErrorResponse(userRequest, CompanyErrors.GeneralError, "Can not create TCR at this moment"); return makeErrorResponse(userRequest, AccountError.GeneralError, "Can not create TCR: " + e.getMessage());
} }
} }
@ -224,13 +219,8 @@ public class ClientCodeService extends QueueConsumer implements InitializingBean
return null; return null;
} }
private RequestInfoUpdate makeErrorResponse(BaseRequest<?> req, CompanyErrors err, Object... arg) { private RequestInfoUpdate makeErrorResponse(BaseRequest<?> req, AccountError err, Object... arg) {
String errMsg = messageResolver.resolve(new EnumMessage(err, arg)); return requestHelper.makeErrorResponse(req, err, arg);
return new RequestInfoUpdate()
.setId(req.getId())
.setStatus(ru.spcex.clearing.platform.messaging.service.Status.Error)
.setMessage(errMsg);
} }
protected RequestInfoUpdate clientCodeDelete(BaseRequest<CommonDeleteRequest> userRequest) { protected RequestInfoUpdate clientCodeDelete(BaseRequest<CommonDeleteRequest> userRequest) {
@ -296,14 +286,15 @@ public class ClientCodeService extends QueueConsumer implements InitializingBean
requestPayload.setTradingClearingRegistryType(TradingClearingRegistryType.Client_B.getKey()); requestPayload.setTradingClearingRegistryType(TradingClearingRegistryType.Client_B.getKey());
request.setRequestPayload(requestPayload); request.setRequestPayload(requestPayload);
try { try {
Long reply = accountServiceExchanger.exchange(request); // Следующий вызываемый метод обязательно должен быть synchronized.
if (reply == null) { RequestInfoUpdate reply = tradingClearingRegistryService.tradingClearingRegistryNew(request);
throw new ClearingBaseException(CompanyErrors.GeneralError, "Waiting account-service timeout"); if (reply != null && Status.Error.equals(reply.getStatus())) {
throw new ClearingBaseException(AccountError.GeneralError, "tradingClearingRegistryService return error: " + reply.getMessage());
} }
// примечание: ошибка и сбой (непредвиденное завершение программы) не приведёт к необратимым последствиям, } catch (ClearingBaseException expected) {
// т.к. пользователь сможет повторить запрос, а созданный на предыдущем запросе ТКР уже будет создан и найдётся. throw expected;
} catch (InterruptedException e) { } catch (Exception e) {
throw new ClearingBaseException(CompanyErrors.GeneralError, "Waiting account-service timeout"); throw new ClearingBaseException(AccountError.GeneralError, "Waiting account-service timeout");
} }
} }

View file

@ -272,7 +272,7 @@ public class TradingClearingRegistryService extends QueueConsumer implements Ini
return null; return null;
} }
public RequestInfoUpdate tradingClearingRegistryNew(BaseRequest<TradingClearingRegistryNewRequest> userRequest) { public synchronized RequestInfoUpdate tradingClearingRegistryNew(BaseRequest<TradingClearingRegistryNewRequest> userRequest) {
log.debug("TradingClearingRegistryNewRequest received"); log.debug("TradingClearingRegistryNewRequest received");
RequestInfoUpdate requestInfoUpdate = userRoleVerification.validateRoleAndGetResult(userRequest); RequestInfoUpdate requestInfoUpdate = userRoleVerification.validateRoleAndGetResult(userRequest);
if (requestInfoUpdate != null) { if (requestInfoUpdate != null) {
@ -328,39 +328,9 @@ public class TradingClearingRegistryService extends QueueConsumer implements Ini
tradingClearingRegistryImdg.insert(tradingClearingRegistry); tradingClearingRegistryImdg.insert(tradingClearingRegistry);
log.debug("successfully processed, id {}", id); log.debug("successfully processed, id {}", id);
sendResponse(userRequest, null, null); // to DESTINATION_TRADING_CLEARING_REGISTRY_REPLY
return null; return null;
} }
private void sendResponse(BaseRequest<?> o, Object response, Header correlationId) {
try {
Future<RecordMetadata> send;
BaseRequest<Object> req = new BaseRequest<>();
req.setId(o.getId());
req.setActionType(ActionType.SYSTEM);
if (response != null) {
req.setRequestPayload(response);
} else {
//default response
RequestInfoUpdate success = new RequestInfoUpdate();
success.setId(o.getId());
success.setStatus(Status.Success);
req.setRequestPayload(success);
}
ProducerRecord<String, Object> respRec = new ProducerRecord<>(Consts.DESTINATION_TRADING_CLEARING_REGISTRY_REPLY, req);
if (correlationId != null) {
respRec.headers().add(new CorrelationHeader(correlationId.value().clone()));
}
send = kafkaProducer.send(respRec);
send.get();
} catch (Exception e) {
log.error(ExceptionUtils.getStackTrace(e));
}
}
public RequestInfoUpdate tradingClearingRegistryUpdate(BaseRequest<TradingClearingRegistryUpdateRequest> userRequest) { public RequestInfoUpdate tradingClearingRegistryUpdate(BaseRequest<TradingClearingRegistryUpdateRequest> userRequest) {
log.debug("TradingClearingRegistryUpdateRequest received"); log.debug("TradingClearingRegistryUpdateRequest received");
RequestInfoUpdate requestInfoUpdate = userRoleVerification.validateRoleAndGetResult(userRequest); RequestInfoUpdate requestInfoUpdate = userRoleVerification.validateRoleAndGetResult(userRequest);

View file

@ -1,4 +1,4 @@
package ru.spcex.clearing.company.service; package ru.spcex.clearing.account.service;
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;
@ -8,58 +8,62 @@ import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith; import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.ArgumentCaptor; import org.mockito.ArgumentCaptor;
import org.mockito.Captor; import org.mockito.Captor;
import org.mockito.Mockito;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.boot.test.mock.mockito.SpyBean; import org.springframework.boot.test.mock.mockito.SpyBean;
import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit.jupiter.SpringExtension; import org.springframework.test.context.junit.jupiter.SpringExtension;
import ru.clearing.classes.statics.data.account.Account; import ru.clearing.classes.statics.data.account.Account;
import ru.clearing.classes.statics.data.account.ClearingAccount;
import ru.clearing.classes.statics.data.account.ClientCode; import ru.clearing.classes.statics.data.account.ClientCode;
import ru.clearing.classes.statics.data.account.DepoAccount;
import ru.clearing.classes.statics.data.company.Company; import ru.clearing.classes.statics.data.company.Company;
import ru.clearing.classes.statics.data.profile.CompanyInfo; import ru.clearing.classes.statics.data.profile.CompanyInfo;
import ru.clearing.classes.statics.data.registry.TradingClearingRegistry; import ru.clearing.classes.statics.data.registry.TradingClearingRegistry;
import ru.clearing.platform.dictionary.*; import ru.clearing.platform.dictionary.*;
import ru.spcex.clearing.company.config.BeanConfiguration; import ru.spcex.clearing.account.config.BeanConfiguration;
import ru.spcex.clearing.company.config.validation.ClientCodeValidationConfig; import ru.spcex.clearing.account.config.HazelcastServiceTestConfiguration;
import ru.spcex.clearing.company.config.validation.ValidationConfig; import ru.spcex.clearing.account.config.KafkaConfigTest;
import ru.spcex.clearing.account.config.validation.ClientCodeValidationConfig;
import ru.spcex.clearing.account.config.validation.TradingClearingRegistryValidationConfig;
import ru.spcex.clearing.account.config.validation.ValidationConfig;
import ru.spcex.clearing.account.utils.MatcherFactory;
import ru.spcex.clearing.account.utils.TestUtils;
import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.imdg.IMDGDistributedNames;
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.account.ClientCodeUpdateRequest; import ru.spcex.clearing.platform.messaging.domain.cud.account.ClientCodeUpdateRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteRequest; import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteRequest;
import ru.spcex.clearing.test.MatcherFactory;
import ru.spcex.clearing.test.TestUtils;
import ru.spcex.clearing.test.config.ImdgTestConfig;
import ru.spcex.clearing.test.config.KafkaTestConfig;
import ru.spcex.platform.enumeration.TradingClearingRegistryType; import ru.spcex.platform.enumeration.TradingClearingRegistryType;
import ru.spcex.platform.enumeration.WorkflowStatus; 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.iml.hazelcast.service.HazelcastService;
import javax.annotation.PostConstruct; import javax.annotation.PostConstruct;
import static org.junit.jupiter.api.Assertions.*; import static org.junit.jupiter.api.Assertions.*;
import static org.mockito.Mockito.*;
import static ru.spcex.clearing.test.MatcherFactory.usingIgnoringFieldsComparator;
import static ru.spcex.clearing.test.TestUtils.*;
import static ru.spcex.clearing.test.config.ImdgTestConfig.waitAvailableImdgProviderAndAddAdminWithDefaultId;
@ExtendWith(SpringExtension.class) @ExtendWith(SpringExtension.class)
@ContextConfiguration(classes = { @ContextConfiguration(classes = {
ClientCodeService.class, ClientCodeService.class,
ClientCodeValidationConfig.class, ClientCodeValidationConfig.class,
TradingClearingRegistryService.class,
TradingClearingRegistryValidationConfig.class,
ValidationConfig.class, ValidationConfig.class,
BeanConfiguration.class, BeanConfiguration.class,
KafkaTestConfig.class, KafkaConfigTest.class,
ImdgTestConfig.class})
HazelcastServiceTestConfiguration.class,})
class ClientCodeServiceTest { class ClientCodeServiceTest {
private static final int PARTITION = 0; private static final int PARTITION = 0;
private static final Long ID = 4L; private static final Long ID = 4L;
public static final MatcherFactory.Matcher<ClientCode> CLIENT_CODE_MATCHER = usingIgnoringFieldsComparator("created","updated"); public static final MatcherFactory.Matcher<ClientCode> CLIENT_CODE_MATCHER = MatcherFactory.usingIgnoringFieldsComparator("created", "updated");
private static final Long TCR_ID = 41L; private static final Long TCR_ID = 41L;
private static final Long COMPANY_ID = 42L; private static final Long COMPANY_ID = 42L;
@ -68,7 +72,7 @@ class ClientCodeServiceTest {
ClientCodeService clientCodeService; ClientCodeService clientCodeService;
@Autowired @Autowired
@Qualifier("hazelcastServiceTest") @Qualifier("hazelcastServiceTest")
private ImdgProvider hazelcastServiceTest; private HazelcastService hazelcastServiceTest;
@Captor @Captor
private ArgumentCaptor<ProducerRecord> producerRecord; private ArgumentCaptor<ProducerRecord> producerRecord;
@SpyBean @SpyBean
@ -81,7 +85,7 @@ class ClientCodeServiceTest {
@PostConstruct @PostConstruct
private void init() { private void init() {
waitAvailableImdgProviderAndAddAdminWithDefaultId(); hazelcastServiceTest.waitAvailable();
clientCodeImdg = hazelcastServiceTest.getImdg(IMDGDistributedNames.Map_ClientCode, ClientCode.class); clientCodeImdg = hazelcastServiceTest.getImdg(IMDGDistributedNames.Map_ClientCode, ClientCode.class);
// Словари для теста, применяются в ValidationConfig // Словари для теста, применяются в ValidationConfig
@ -120,21 +124,33 @@ class ClientCodeServiceTest {
tcrImdg.insert(registry1); tcrImdg.insert(registry1);
Imdg<Account> accounts = hazelcastServiceTest.getImdg(IMDGDistributedNames.Map_Account, Account.class); Imdg<Account> accounts = hazelcastServiceTest.getImdg(IMDGDistributedNames.Map_Account, Account.class);
Account moneyAccount=new Account(); Imdg<ClearingAccount> clsAccounts = hazelcastServiceTest.getImdg(IMDGDistributedNames.Map_ClearingAccount, ClearingAccount.class);
Account moneyAccount = new Account();
moneyAccount.setId(131L); moneyAccount.setId(131L);
moneyAccount.setAccount("AAAA-4444"); moneyAccount.setAccount("AAAA-4444");
moneyAccount.setStatus("ACTV"); moneyAccount.setStatus("ACTV");
moneyAccount.setCompanyId(COMPANY_ID); // для валидации принадлежности счёта moneyAccount.setCompanyId(COMPANY_ID); // для валидации принадлежности счёта
accounts.insert(moneyAccount); accounts.insert(moneyAccount);
Account depoAccount=new Account(); ClearingAccount clsAcc = new ClearingAccount();
clsAcc.setId(moneyAccount.getId());
clsAcc.setCompanyId(COMPANY_ID);
clsAcc.setAccountId(moneyAccount.getId());
clsAccounts.insert(clsAcc);
Imdg<DepoAccount> depoAccounts = hazelcastServiceTest.getImdg(IMDGDistributedNames.Map_DepoAccount, DepoAccount.class);
Account depoAccount = new Account();
depoAccount.setId(132L); depoAccount.setId(132L);
depoAccount.setAccount("AAAB-44654"); depoAccount.setAccount("AAAB-44654");
depoAccount.setStatus("ACTV"); depoAccount.setStatus("ACTV");
depoAccount.setCompanyId(COMPANY_ID); // для валидации принадлежности счёта depoAccount.setCompanyId(COMPANY_ID); // для валидации принадлежности счёта
accounts.insert(depoAccount); accounts.insert(depoAccount);
DepoAccount depoAcc = new DepoAccount();
depoAcc.setId(depoAccount.getId());
depoAcc.setCompanyId(COMPANY_ID);
depoAcc.setAccountId(depoAccount.getId());
depoAccounts.insert(depoAcc);
TestUtils.FutureRecordMetadata future = spy(TestUtils.FutureRecordMetadata.class); TestUtils.FutureRecordMetadata future = Mockito.spy(TestUtils.FutureRecordMetadata.class);
doReturn(future).when(mockProducer).send(producerRecord.capture()); Mockito.doReturn(future).when(mockProducer).send(producerRecord.capture());
} }
private <D extends AbstractDictionary> void putToDictionary(String mapName, D object, String code) { private <D extends AbstractDictionary> void putToDictionary(String mapName, D object, String code) {
@ -187,12 +203,12 @@ class ClientCodeServiceTest {
// predictableClientCode.setTradingClearingRegistryId(TCR_ID); // predictableClientCode.setTradingClearingRegistryId(TCR_ID);
//ACT //ACT
String jsonString = getJsonStringForNew(clientCodeNewRequest, ID); String jsonString = TestUtils.getJsonStringForNew(clientCodeNewRequest, ID);
addRecordToKafka((MockConsumer) clientCodeService.getConsumer(), Consts.DESTINATION_CLIENT_CODE_NEW, PARTITION, 0, jsonString); TestUtils.addRecordToKafka((MockConsumer) clientCodeService.getConsumer(), Consts.DESTINATION_CLIENT_CODE_NEW, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); TestUtils.waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord);
ClientCode resultNew = clientCodeImdg.getSingleObjectBySQL(String.format("code = '%s'", ccCode)); ClientCode resultNew = clientCodeImdg.getSingleObjectBySQL(String.format("code = '%s'", ccCode));
predictableClientCode.setId(resultNew.getId()); predictableClientCode.setId(resultNew.getId());
CLIENT_CODE_MATCHER.assertMatch(resultNew, predictableClientCode); CLIENT_CODE_MATCHER.assertMatch(resultNew, predictableClientCode);
@ -225,19 +241,57 @@ class ClientCodeServiceTest {
// predictableClientCode.setTradingClearingRegistryId(TCR_ID); // predictableClientCode.setTradingClearingRegistryId(TCR_ID);
//ACT //ACT
String jsonString = getJsonStringForNew(clientCodeNewRequest, ID); String jsonString = TestUtils.getJsonStringForNew(clientCodeNewRequest, ID);
addRecordToKafka((MockConsumer) clientCodeService.getConsumer(), Consts.DESTINATION_CLIENT_CODE_NEW_UM_COMPANY, PARTITION, 0, jsonString); TestUtils.addRecordToKafka((MockConsumer) clientCodeService.getConsumer(), Consts.DESTINATION_CLIENT_CODE_NEW_UM_COMPANY, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); TestUtils.waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord);
ClientCode resultNew = clientCodeImdg.getSingleObjectBySQL(String.format("code = '%s'", ccCode)); ClientCode resultNew = clientCodeImdg.getSingleObjectBySQL(String.format("code = '%s'", ccCode));
predictableClientCode.setId(resultNew.getId()); predictableClientCode.setId(resultNew.getId());
CLIENT_CODE_MATCHER.assertMatch(resultNew, predictableClientCode); CLIENT_CODE_MATCHER.assertMatch(resultNew, predictableClientCode);
assertNotNull(resultNew.getCreated()); assertNotNull(resultNew.getCreated());
} }
//todo добавить тест NEW заполненными MoneyAccountId(131L), DepoAccountId(132L); - от этого направляется дополнительное сообщение в очередь и используется ожидание ответа.
/**
* {@link ClientCodeService#clientCodeUpdate(BaseRequest)}<br>
* Тест проверяет создание {@link ClientCode} в IMDG при передаче из Apache Kafka (очередь 1).<br>
* NEW с заполненными MoneyAccountId(131L), DepoAccountId(132L); - от этого будет создан ТКР.
* Входной запрос {@link ru.spcex.clearing.platform.messaging.domain.cud.account.ClientCodeNewRequest}:<br>
**/
@Test
void clientCodeNew3() {
//ARRANGE
final String ccCode = "Lucky planet2";
ClientCodeNewRequest clientCodeNewRequest = new ClientCodeNewRequest();
clientCodeNewRequest.setCompanyId(COMPANY_ID);
clientCodeNewRequest.setCode(ccCode);
clientCodeNewRequest.setTradingClearingRegistryId(null);
clientCodeNewRequest.setMoneyAccountId(131L);
clientCodeNewRequest.setDepoAccountId(132L);
clientCodeNewRequest.setStatus("ACTV");
ClientCode predictableClientCode = new ClientCode();
predictableClientCode.setCode(ccCode);
predictableClientCode.setStatus("ACTV");
predictableClientCode.setCompanyId(COMPANY_ID);
predictableClientCode.setMoneyAccountId(131L);
predictableClientCode.setDepoAccountId(132L);
predictableClientCode.setTradingClearingRegistryId(TCR_ID);
//ACT
String jsonString = TestUtils.getJsonStringForNew(clientCodeNewRequest, ID);
TestUtils.addRecordToKafka((MockConsumer) clientCodeService.getConsumer(), Consts.DESTINATION_CLIENT_CODE_NEW, PARTITION, 0, jsonString);
//ASSERT
TestUtils.waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord);
ClientCode resultNew = clientCodeImdg.getSingleObjectBySQL(String.format("code = '%s'", ccCode));
predictableClientCode.setId(resultNew.getId());
CLIENT_CODE_MATCHER.assertMatch(resultNew, predictableClientCode);
assertNotNull(resultNew.getCreated());
}
/** /**
* {@link ClientCodeService#clientCodeUpdate(BaseRequest)}<br> * {@link ClientCodeService#clientCodeUpdate(BaseRequest)}<br>
@ -276,12 +330,12 @@ class ClientCodeServiceTest {
predictableClientCode.setStatus("ACTV"); predictableClientCode.setStatus("ACTV");
//ACT //ACT
String jsonString = getJsonStringForUpdate(clientCodeUpdateRequest, ID); String jsonString = TestUtils.getJsonStringForUPDATE(clientCodeUpdateRequest, ID);
addRecordToKafka((MockConsumer) clientCodeService.getConsumer(), Consts.DESTINATION_CLIENT_CODE_UPDATE, PARTITION, 0, jsonString); TestUtils.addRecordToKafka((MockConsumer) clientCodeService.getConsumer(), Consts.DESTINATION_CLIENT_CODE_UPDATE, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); TestUtils.waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord);
ClientCode resultUpdating = clientCodeImdg.getSingleObjectByID(ID); ClientCode resultUpdating = clientCodeImdg.getSingleObjectByID(ID);
CLIENT_CODE_MATCHER.assertMatch(resultUpdating, predictableClientCode); CLIENT_CODE_MATCHER.assertMatch(resultUpdating, predictableClientCode);
@ -315,13 +369,13 @@ class ClientCodeServiceTest {
Assertions.assertNotNull(clientCodeImdg.getSingleObjectByID(ID)); // verify test data Assertions.assertNotNull(clientCodeImdg.getSingleObjectByID(ID)); // verify test data
//ACT //ACT
String jsonString = getJsonStringForUpdate(clientCodeDeleteRequest, ID); String jsonString = TestUtils.getJsonStringForUPDATE(clientCodeDeleteRequest, ID);
addRecordToKafka((MockConsumer) clientCodeService.getConsumer(), Consts.DESTINATION_CLIENT_CODE_DELETE, PARTITION, 0, jsonString); TestUtils.addRecordToKafka((MockConsumer) clientCodeService.getConsumer(), Consts.DESTINATION_CLIENT_CODE_DELETE, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); TestUtils.waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord);
ClientCode resultUpdate = clientCodeImdg.getSingleObjectByID(ID); ClientCode resultUpdate = clientCodeImdg.getSingleObjectByID(ID);
Assertions.assertNull(resultUpdate); Assertions.assertNull(resultUpdate);
} }

File diff suppressed because it is too large Load diff

View file

@ -89,7 +89,6 @@ public interface Consts {
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"; // todo check buplicate 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_REPLY = "trading-clearing-registry-reply";
String DESTINATION_TRADING_CLEARING_REGISTRY_BLOCK = "trading-clearing-registry-block"; String DESTINATION_TRADING_CLEARING_REGISTRY_BLOCK = "trading-clearing-registry-block";
String DESTINATION_SDF08_NEW = "s-df-08-new"; String DESTINATION_SDF08_NEW = "s-df-08-new";

View file

@ -163,7 +163,7 @@ public class QueueConsumer implements AutoCloseable {
}); });
} }
private void sendResponse(BaseRequest<?> o, Object response, Header correlationId) { protected void sendResponse(BaseRequest<?> o, Object response, Header correlationId) {
outputExecutor.submit(() -> { outputExecutor.submit(() -> {
try { try {
Future<RecordMetadata> send; Future<RecordMetadata> send;