This commit is contained in:
AKurakin 2023-05-04 16:17:21 +03:00
commit f1a803aacf
8 changed files with 1156 additions and 19 deletions

View file

@ -1,25 +1,56 @@
package ru.spcex.clearing.util.services.exchangers;
import org.apache.commons.lang3.StringUtils;
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.clients.producer.RecordMetadata;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.DisposableBean;
import org.springframework.beans.factory.InitializingBean;
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.balance.AccountBalanceClearingRequest;
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.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.platform.enumeration.Task;
import ru.spcex.platform.utils.log.ExceptionUtils;
import java.io.Closeable;
import java.util.Objects;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.Future;
import java.util.concurrent.atomic.AtomicInteger;
/**
* Синхронный обмен сообщениями с ассинхронным сервисом.
* Метод exchange() отправляет сообщение в очередь и дожидается ответа из другой очереди
* Метод exchange() отправляет сообщение в очередь и дожидается ответа из другой очереди.
* Использовать с осторожностью, т.к. нет возобновления ожидания в случае нештатного выключения модуля при ожидании ответа из kafka.
* <p>
* <p>
* Ожидает из очереди (CommonIdRequest)
*
* @param <TOut> отправляется в очередь
*/
public class BiDirectionQueueExchanger<TIn, TOut extends BaseRequest<?>> extends QueueConsumer implements InitializingBean, DisposableBean {
public class BiDirectionQueueExchanger<TOut extends BaseRequest<?>> extends QueueConsumer implements Closeable {
protected final Logger log = LoggerFactory.getLogger(getClass());
protected final Object sync = new Object();
protected String outQueue;
protected String inQueue;
protected Class<TIn> listenClass;
protected long timeout;
private Producer<String, Object> producer;
protected volatile boolean isTerminated;
protected Object syncObject;
protected Long lastSentRequestId;
protected volatile Long lastResponseId;
protected boolean ignoreOtherResponse = true; // Пропускать другие ID, пока не получит lastResponseId==lastSentRequestId
/**
* Синхронно-ассинхронный обмен сообщениями
@ -27,40 +58,159 @@ public class BiDirectionQueueExchanger<TIn, TOut extends BaseRequest<?>> extends
* @param kafkaQueue
* @param kafkaProducer
* @param outQueue отправляет в очередь
* @param inQueue слушает очередь, ожидает ответов
* @param listenClass типы объектов из inQueue
* @param timeout - максимальное ожидание ответа, в миллисекундах, 0 - неограничено
* @param inQueue слушает очередь/топик, ожидает ответов. Пример: Consts.CONTINUE_CLEARING
* @param ignoreOtherResponse true для inQueue=Consts.REQUEST_INFO_UPDATE
* @param timeout - максимальное ожидание ответа, в миллисекундах, 0 - неограничено
*/
public BiDirectionQueueExchanger(Consumer<String, Object> kafkaQueue, Producer<String, Object> kafkaProducer,
String outQueue,
String inQueue, Class<TIn> listenClass,
String inQueue,
boolean ignoreOtherResponse,
long timeout) {
super(kafkaQueue, kafkaProducer);
super(kafkaQueue);
this.producer = kafkaProducer;
if (kafkaQueue == null) {
throw new IllegalArgumentException("kafkaQueue is empty");
}
if (kafkaProducer == null) {
throw new IllegalArgumentException("kafkaProducer is empty");
}
if (StringUtils.isEmpty(outQueue)) {
throw new IllegalArgumentException("outQueue is empty");
}
this.outQueue = outQueue;
if (StringUtils.isEmpty(inQueue)) {
throw new IllegalArgumentException("inQueue is empty");
}
this.inQueue = inQueue;
this.listenClass = listenClass;
if (timeout < 0) {
timeout = 0;
}
this.timeout = timeout;
this.ignoreOtherResponse = ignoreOtherResponse;
initReplyListener();
}
/**
* Отправить сообщение message в outQueue и дождаться ответа из очереди inQueue
*
* @param message
* @return
* @throws InterruptedException
* @return request/reply id
* @throws InterruptedException или когда поток прерывают, или когда BiDirectionQueueExchanger.close()
*/
public TOut exchange(TIn message) throws InterruptedException {
//todo impl BiDirectionQueueExchanger
return null;
public Long exchange(TOut message) throws InterruptedException {
if (syncObject == null)
throw new IllegalStateException("Listener queue not initialized");
if (message.getId() == null) {
throw new IllegalArgumentException("Required message " + message.getClass().getSimpleName() + ".id is null");
}
//send request to kafka, wait for a reply
try {
synchronized (syncObject) {
Long sentRequestId = sendMessage(message);
log.trace("sent request to kafka, requestId: {}", sentRequestId);
this.lastSentRequestId = sentRequestId;
long waitStart = System.currentTimeMillis();
long waitTime = timeout;
do {
try {
if (timeout > 0) {
syncObject.wait(waitTime);
} else {
syncObject.wait();
}
//если на этом месте произошла ошибка - непонятно как обрабатывать
//т.к. request на самом деле могла быть проблема с сетью/недоступностью кафки etc.
} catch (InterruptedException e) {
log.error(ExceptionUtils.getStackTrace(e));
throw e;
// Thread.currentThread().interrupt();
// throw new RuntimeException(e);
}
waitTime = timeout - (System.currentTimeMillis() - waitStart);
if (isTerminated && lastResponseId == null) {
log.warn("Terminated waiter {} kafka at thread {}", this, Thread.currentThread().getName());
throw new InterruptedException(getClass().getSimpleName() + " closed");
}
if (ignoreOtherResponse && !sentRequestId.equals(lastResponseId)) {
log.debug("Ignore response id={}, we waiting {}.", lastResponseId, sentRequestId);
}
} while (ignoreOtherResponse && !sentRequestId.equals(lastResponseId) && (timeout > 0 && waitTime > 0));
if (!sentRequestId.equals(lastResponseId)) {
log.error("FATAL: last response from {} update from Kafka ID was {}, but sent ID was {}",
inQueue, lastResponseId, sentRequestId);
throw new IllegalStateException("Expected wait till requestId=" + sentRequestId + " but catch lastResponseId=" + lastResponseId);
}
return sentRequestId;
}
} finally {
this.lastSentRequestId = null;
}
}
@Override
public void destroy() throws Exception {
public void close() {
isTerminated = true;
if (lastSentRequestId != null) {
log.warn("Finished before wait, lastSentRequestId={}", lastSentRequestId);
if (syncObject != null) {
synchronized (syncObject) {
lastResponseId = null;
syncObject.notifyAll();
}
}
}
syncObject = null;
log.debug("Listener {} for async exchange {}-{} close", this, outQueue, inQueue);
}
@Override
public void afterPropertiesSet() throws Exception {
protected void initReplyListener() {
if (syncObject != null) {
throw new IllegalStateException("Already initialized");
}
callback(CommonIdRequest.class)
.setConsumer(this::continueWaiting)
.forDestination(inQueue, callbacks::put);
init();
syncObject = new Object();
isTerminated = false;
log.debug("Listener {} for async exchange {}-{} ready", this, outQueue, inQueue);
}
protected Long sendMessage(TOut message) throws InterruptedException {
Long sentRequestId = message.getId();
if (sentRequestId == null) {
throw new IllegalArgumentException("Message " + message.getClass().getSimpleName() + " required id.");
}
try {
Future<RecordMetadata> send = producer.send(new ProducerRecord<>(Consts.REQUEST_INFO_UPDATE, message));
send.get();
return sentRequestId;
} catch (InterruptedException e) {
log.error("Interrupt when send message, {}", ExceptionUtils.getStackTrace(e));
Thread.currentThread().interrupt();
throw e;
} catch (ExecutionException e) {
log.error("Error at send message, {} cause: {}", e.toString(), ExceptionUtils.getStackTrace(e.getCause()));
throw new RuntimeException("Send message error", e.getCause());
} catch (Exception e) {
log.error("Unexpected exception when send message: {}", ExceptionUtils.getStackTrace(e));
}
return null;
}
protected void continueWaiting(BaseRequest<CommonIdRequest> event) {
synchronized (syncObject) {
CommonIdRequest requestPayload = event.getRequestPayload();
log.debug("continueWaiting {} for id={}", inQueue, requestPayload.getId());
lastResponseId = requestPayload.getId();
syncObject.notifyAll();
}
}
}

View file

@ -0,0 +1,182 @@
package ru.spcex.clearing.company.config.validation;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import ru.clearing.classes.statics.data.account.Account;
import ru.clearing.classes.statics.data.account.ClientCode;
import ru.clearing.classes.statics.data.company.Company;
import ru.clearing.classes.statics.data.registry.TradingClearingRegistry;
import ru.clearing.platform.dictionary.WorkflowStatusDictionary;
import ru.spcex.clearing.company.error.CompanyErrors;
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.ClientCodeUpdateRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteRequest;
import ru.spcex.clearing.validation.common.rules.DictionaryPresentRule;
import ru.spcex.clearing.validation.common.rules.FieldRequiredRule;
import ru.spcex.clearing.validation.common.rules.IdPresentRule;
import ru.spcex.platform.classes.base.SpcexObjectBase;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.validation.ImdgValidationContext;
import ru.spcex.platform.utils.validation.IValidator;
import ru.spcex.platform.utils.validation.ValidatorImpl;
import java.util.Map;
import java.util.Objects;
import java.util.function.Consumer;
import java.util.function.Function;
@Configuration
public class ClientCodeValidationConfig {
@Bean("clientCodeNewRequestValidator")
public Function<ClientCodeNewRequest, IValidator> clientCodeNewRequestValidator(Map<String, Imdg<? extends SpcexObjectBase>> imdgForValidation) {
return clientCodeUpdateRequest -> {
ImdgValidationContext<ClientCodeNewRequest> context = new ImdgValidationContext<>();
context.setValidatedObject(clientCodeUpdateRequest);
Consumer<String> addImdg = (s) -> context.addImdg(s, imdgForValidation.get(s));
addImdg.accept(IMDGDistributedNames.Map_Company);
addImdg.accept(IMDGDistributedNames.Map_Account);
addImdg.accept(IMDGDistributedNames.Map_TradingClearingRegistry);
addImdg.accept(IMDGDistributedNames.Map_WorkflowStatusDictionary);
return new ValidatorImpl<>(context,
FieldRequiredRule.instance("companyId", ClientCodeNewRequest::getCompanyId, CompanyErrors.RequiredFieldEmpty),
IdPresentRule.instance("companyId",
ClientCodeNewRequest::getCompanyId,
IMDGDistributedNames.Map_Company,
Company.class,
CompanyErrors.RequiredFieldEmpty,
CompanyErrors.CompanyNotFound),
IdPresentRule.instance("moneyAccountId",
ClientCodeNewRequest::getMoneyAccountId,
IMDGDistributedNames.Map_Account,
Account.class,
CompanyErrors.RequiredFieldEmpty,
CompanyErrors.AccountNotFound,
false,
acc -> Objects.equals(context.getValidatedObject().getCompanyId(), acc.getCompanyId())
? null : CompanyErrors.AccountNotFound
),
IdPresentRule.instance("depoAccountId",
ClientCodeNewRequest::getDepoAccountId,
IMDGDistributedNames.Map_Account,
Account.class,
CompanyErrors.RequiredFieldEmpty,
CompanyErrors.AccountNotFound,
false,
acc -> Objects.equals(context.getValidatedObject().getCompanyId(), acc.getCompanyId())
? null : CompanyErrors.AccountNotFound
),
IdPresentRule.instance("tradingClearingRegistryId",
ClientCodeNewRequest::getTradingClearingRegistryId,
IMDGDistributedNames.Map_TradingClearingRegistry,
TradingClearingRegistry.class,
CompanyErrors.RequiredFieldEmpty,
CompanyErrors.TradingClearingRegistryNotFound,
false,
acc -> Objects.equals(context.getValidatedObject().getCompanyId(), acc.getCompanyId())
? null : CompanyErrors.TradingClearingRegistryNotFound
),
DictionaryPresentRule.instance("stauts",
ClientCodeNewRequest::getStatus,
IMDGDistributedNames.Map_WorkflowStatusDictionary,
WorkflowStatusDictionary.class,
CompanyErrors.RequiredFieldEmpty,
CompanyErrors.DictionaryNotFound,
false)
);
};
}
@Bean("clientCodeUpdateRequestValidator")
public Function<ClientCodeUpdateRequest, IValidator> clientCodeUpdateRequestValidator(Map<String, Imdg<? extends SpcexObjectBase>> imdgForValidation) {
return clientCodeUpdateRequest -> {
ImdgValidationContext<ClientCodeUpdateRequest> context = new ImdgValidationContext<>();
context.setValidatedObject(clientCodeUpdateRequest);
Consumer<String> addImdg = (s) -> context.addImdg(s, imdgForValidation.get(s));
addImdg.accept(IMDGDistributedNames.Map_ClientCode);
addImdg.accept(IMDGDistributedNames.Map_Company);
addImdg.accept(IMDGDistributedNames.Map_Account);
addImdg.accept(IMDGDistributedNames.Map_TradingClearingRegistry);
addImdg.accept(IMDGDistributedNames.Map_WorkflowStatusDictionary);
return new ValidatorImpl<ImdgValidationContext<ClientCodeUpdateRequest>>(context,
IdPresentRule.instance("id",
ClientCodeUpdateRequest::getId,
IMDGDistributedNames.Map_ClientCode,
ClientCode.class,
CompanyErrors.RequiredFieldEmpty,
CompanyErrors.RecordNotFound
),
FieldRequiredRule.instance("companyId", ClientCodeUpdateRequest::getCompanyId, CompanyErrors.RequiredFieldEmpty),
IdPresentRule.instance("companyId",
ClientCodeUpdateRequest::getCompanyId,
IMDGDistributedNames.Map_Company,
Company.class,
CompanyErrors.RequiredFieldEmpty,
CompanyErrors.CompanyNotFound),
IdPresentRule.instance("moneyAccountId",
ClientCodeUpdateRequest::getMoneyAccountId,
IMDGDistributedNames.Map_Account,
Account.class,
CompanyErrors.RequiredFieldEmpty,
CompanyErrors.AccountNotFound,
false,
acc -> Objects.equals(context.getValidatedObject().getCompanyId(), acc.getCompanyId())
? null : CompanyErrors.AccountNotFound
),
IdPresentRule.instance("depoAccountId",
ClientCodeUpdateRequest::getDepoAccountId,
IMDGDistributedNames.Map_Account,
Account.class,
CompanyErrors.RequiredFieldEmpty,
CompanyErrors.AccountNotFound,
false,
acc -> Objects.equals(context.getValidatedObject().getCompanyId(), acc.getCompanyId())
? null : CompanyErrors.AccountNotFound
),
IdPresentRule.instance("tradingClearingRegistryId",
ClientCodeUpdateRequest::getTradingClearingRegistryId,
IMDGDistributedNames.Map_TradingClearingRegistry,
TradingClearingRegistry.class,
CompanyErrors.RequiredFieldEmpty,
CompanyErrors.TradingClearingRegistryNotFound,
false,
acc -> Objects.equals(context.getValidatedObject().getCompanyId(), acc.getCompanyId())
? null : CompanyErrors.TradingClearingRegistryNotFound
),
DictionaryPresentRule.instance("stauts",
ClientCodeUpdateRequest::getStatus,
IMDGDistributedNames.Map_WorkflowStatusDictionary,
WorkflowStatusDictionary.class,
CompanyErrors.RequiredFieldEmpty,
CompanyErrors.DictionaryNotFound,
false)
);
};
}
@Bean("clientCodeDeleteRequestValidator")
public Function<CommonDeleteRequest, IValidator> clientCodeDeleteRequestValidator(Map<String, Imdg<? extends SpcexObjectBase>> imdgForValidation) {
return companyDeleteRequest -> {
ImdgValidationContext<CommonDeleteRequest> context = new ImdgValidationContext<>();
context.setValidatedObject(companyDeleteRequest);
Consumer<String> addImdg = (s) -> context.addImdg(s, imdgForValidation.get(s));
addImdg.accept(IMDGDistributedNames.Map_ClientCode);
return new ValidatorImpl<>(context,
// FieldRequiredRule.instance("id", CommonDeleteRequest::getId, CompanyErrors.RequiredFieldEmpty),
IdPresentRule.instance("id",
CommonDeleteRequest::getId,
IMDGDistributedNames.Map_ClientCode,
ClientCode.class,
CompanyErrors.RequiredFieldEmpty,
CompanyErrors.RecordNotFound)
);
};
}
}

View file

@ -2,11 +2,14 @@ package ru.spcex.clearing.company.config.validation;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import ru.clearing.classes.statics.data.account.Account;
import ru.clearing.classes.statics.data.account.ClientCode;
import ru.clearing.classes.statics.data.company.Company;
import ru.clearing.classes.statics.data.company.CompanySymbols;
import ru.clearing.classes.statics.data.company.ClearingMemberCategory;
import ru.clearing.classes.statics.data.profile.Contact;
import ru.clearing.classes.statics.data.profile.ProfileDocument;
import ru.clearing.classes.statics.data.registry.TradingClearingRegistry;
import ru.clearing.platform.dictionary.*;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.validation.common.ValidationHelper;
@ -44,6 +47,10 @@ public class ValidationConfig {
addImdg.accept(IMDGDistributedNames.Map_LegalKindDictionary, LegalKindDictionary.class);
addImdg.accept(IMDGDistributedNames.Map_OrganizationTypeDictionary, OrganizationTypeDictionary.class);
addImdg.accept(IMDGDistributedNames.Map_ClientCode, ClientCode.class);
addImdg.accept(IMDGDistributedNames.Map_Account, Account.class);
addImdg.accept(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class);
return imdg;
}

View file

@ -22,7 +22,9 @@ public enum CompanyErrors implements IErrorEnumId {
CompanyAlreadyHasLiabilities(3018L), // У компании %s присутствуют обязательства.
CompanyWithCompanySymbolAlreadyExist(3019L), // Компания с %s = %s уже создана" (где первый %s - companySymbol, второй %s - companySymbolValue)
EditCompanySymbols(3020L), // Тип реквизита компании %s не может быть изменен.
EditContactType(3021L) // Тип контакта компании %s не может быть изменен.
EditContactType(3021L), // Тип контакта компании %s не может быть изменен.
TradingClearingRegistryNotFound(3022L), // ТКР с %s не найден
AccountNotFound(3023L), // Счет %s не найден
;
private final Long id;

View file

@ -0,0 +1,416 @@
package ru.spcex.clearing.company.service;
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.clients.producer.RecordMetadata;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.DisposableBean;
import org.springframework.beans.factory.InitializingBean;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.account.ClientCode;
import ru.clearing.classes.statics.data.company.Company;
import ru.clearing.classes.statics.data.registry.TradingClearingRegistry;
import ru.spcex.clearing.company.error.CompanyErrors;
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.account.ClientCodeNewRequest;
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.registry.TradingClearingRegistryNewRequest;
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.RequestInfoUpdate;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.clearing.util.security.UserRoleVerification;
import ru.spcex.clearing.util.services.exchangers.BiDirectionQueueExchanger;
import ru.spcex.clearing.validation.common.ValidationHelper;
import ru.spcex.platform.enumeration.TradingClearingRegistryType;
import ru.spcex.platform.enumeration.WorkflowStatus;
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.imdg.iml.hazelcast.adapter.ImdgIdGeneratorHazelcast;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.enumeration.IMessageResolver;
import ru.spcex.platform.utils.error.ClearingBaseException;
import ru.spcex.platform.utils.log.ExceptionUtils;
import ru.spcex.platform.utils.validation.IValidator;
import java.time.Instant;
import java.util.Collection;
import java.util.HashMap;
import java.util.Map;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.Future;
import java.util.function.Function;
@Service
public class ClientCodeService extends QueueConsumer implements InitializingBean, DisposableBean {
private final Logger log = LoggerFactory.getLogger(getClass());
private final Producer<String, Object> kafkaProducer;
private final ImdgId idGenerator;
private final Imdg<ClientCode> clientCodeMap;
private final Imdg<Company> companyMap;
private final Imdg<TradingClearingRegistry> tradingClearingRegistryMap;
private final Function<ClientCodeNewRequest, IValidator> clientCodeNewRequestValidator;
private final Function<ClientCodeUpdateRequest, IValidator> clientCodeUpdateRequestValidator;
private final Function<CommonDeleteRequest, IValidator> clientCodeDeleteRequestValidator;
private final ValidationHelper validationHelper;
private final UserRoleVerification userRoleVerification;
private final IMessageResolver messageResolver;
BiDirectionQueueExchanger<BaseRequest<TradingClearingRegistryNewRequest>> accountServiceExchanger;
@Autowired
public ClientCodeService(Consumer<String, Object> kafkaQueue,
Producer<String, Object> kafkaProducer,
ImdgProvider imdgProvider,
ValidationHelper validationHelper,
IMessageResolver messageResolver,
UserRoleVerification userRoleVerification,
@Qualifier("clientCodeNewRequestValidator") Function<ClientCodeNewRequest, IValidator> clientCodeNewRequestValidator,
@Qualifier("clientCodeUpdateRequestValidator") Function<ClientCodeUpdateRequest, IValidator> clientCodeUpdateRequestValidator,
@Qualifier("clientCodeDeleteRequestValidator") Function<CommonDeleteRequest, IValidator> clientCodeDeleteRequestValidator) {
super(kafkaQueue, kafkaProducer);
this.kafkaProducer = kafkaProducer;
this.idGenerator = imdgProvider.getImdgIdGenerator();
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.clientCodeNewRequestValidator = clientCodeNewRequestValidator;
this.clientCodeUpdateRequestValidator = clientCodeUpdateRequestValidator;
this.clientCodeDeleteRequestValidator = clientCodeDeleteRequestValidator;
this.validationHelper = validationHelper;
this.messageResolver = messageResolver;
this.userRoleVerification = userRoleVerification;
this.accountServiceExchanger = new BiDirectionQueueExchanger<>(kafkaQueue, kafkaProducer,
Consts.DESTINATION_TRADING_CLEARING_REGISTRY_NEW,
Consts.REQUEST_INFO_UPDATE, true, // REQUEST_INFO_UPDATE - стандартная очередь, для результатов всех реквестов. DESTINATION_TRADING_CLEARING_REGISTRY_REPLY,
60000
);
}
@Override
public void afterPropertiesSet() {
callback(ClientCodeNewRequest.class)
.setConsumer(this::clientCodeNew)
.forDestination(Consts.DESTINATION_CLIENT_CODE_NEW, callbacks::put);
callback(ClientCodeNewRequest.class)
.setConsumer(this::clientCodeNewFromApiUmCompany)
.forDestination(Consts.DESTINATION_CLIENT_CODE_NEW_UM_COMPANY, callbacks::put);
callback(ClientCodeUpdateRequest.class)
.setConsumer(this::clientCodeUpdate)
.forDestination(Consts.DESTINATION_CLIENT_CODE_UPDATE, callbacks::put);
callback(CommonDeleteRequest.class)
.setConsumer(this::clientCodeDelete)
.forDestination(Consts.DESTINATION_CLIENT_CODE_DELETE, callbacks::put);
init();
}
@Override
public void destroy() {
accountServiceExchanger.close();
}
protected RequestInfoUpdate clientCodeNew(BaseRequest<ClientCodeNewRequest> userRequest) {
log.debug("ClientCodeNewRequest received {}", userRequest.getId());
RequestInfoUpdate requestInfoUpdate = userRoleVerification.validateRoleAndGetResult(userRequest);
if (requestInfoUpdate != null) return requestInfoUpdate;
requestInfoUpdate = validationHelper.validateTillFirstError(userRequest, clientCodeNewRequestValidator);
if (requestInfoUpdate != null) return requestInfoUpdate;
ClientCodeNewRequest req = userRequest.getRequestPayload();
boolean doCreateTCR = checkNeedCreateTCR(req.getCompanyId(), req.getMoneyAccountId(), req.getDepoAccountId());
if (doCreateTCR) {
try {
createAndWaitTCR(userRequest.getId(), null,
req.getCompanyId(), req.getMoneyAccountId(), req.getDepoAccountId());
} catch (ClearingBaseException e) {
log.error("Can not wait creation of TCR. request id={};CompanyId={}, MoneyAccountId={}, DepoAccountId={}; {}",
userRequest.getId(), req.getCompanyId(), req.getMoneyAccountId(), req.getDepoAccountId(),
e.toString());
return makeErrorResponse(userRequest, CompanyErrors.GeneralError, "Can not create TCR at this moment");
}
}
if (req.getTradingClearingRegistryId() == null && req.getMoneyAccountId() != null) {
TradingClearingRegistry tradingClearingRegistry = selectTradingClearingRegistry(req.getCompanyId(), req.getMoneyAccountId(), req.getDepoAccountId());
if (tradingClearingRegistry == null) {
return makeErrorResponse(userRequest, CompanyErrors.TradingClearingRegistryNotFound);
} else {
req.setTradingClearingRegistryId(tradingClearingRegistry.getId());
}
}
ClientCode newClientCode = buildClientCode(req);
clientCodeMap.insert(newClientCode);
log.debug("successfully processed, new clientCode id {}", newClientCode.getId());
return null;
}
protected RequestInfoUpdate clientCodeNewFromApiUmCompany(BaseRequest<ClientCodeNewRequest> userRequest) {
log.debug("ClientCodeNewRequest received {}", userRequest.getId());
RequestInfoUpdate requestInfoUpdate = userRoleVerification.validateRoleAndGetResult(userRequest);
if (requestInfoUpdate != null) return requestInfoUpdate;
requestInfoUpdate = validationHelper.validateTillFirstError(userRequest, clientCodeNewRequestValidator);
if (requestInfoUpdate != null) return requestInfoUpdate;
ClientCodeNewRequest req = userRequest.getRequestPayload();
boolean doCreateTCR = checkNeedCreateTCR(req.getCompanyId(), req.getMoneyAccountId(), req.getDepoAccountId());
if (doCreateTCR) {
try {
createAndWaitTCR(userRequest.getId(), null,
req.getCompanyId(), req.getMoneyAccountId(), req.getDepoAccountId());
} catch (ClearingBaseException e) {
log.error("Can not wait creation of TCR. request id={};CompanyId={}, MoneyAccountId={}, DepoAccountId={}; {}",
userRequest.getId(), req.getCompanyId(), req.getMoneyAccountId(), req.getDepoAccountId(),
e.toString());
return makeErrorResponse(userRequest, CompanyErrors.GeneralError, "Can not create TCR at this moment");
}
}
if (req.getTradingClearingRegistryId() == null && req.getMoneyAccountId() != null) {
TradingClearingRegistry tradingClearingRegistry = selectTradingClearingRegistry(req.getCompanyId(), req.getMoneyAccountId(), req.getDepoAccountId());
if (tradingClearingRegistry == null) {
return makeErrorResponse(userRequest, CompanyErrors.TradingClearingRegistryNotFound);
} else {
req.setTradingClearingRegistryId(tradingClearingRegistry.getId());
}
}
ClientCode newClientCode = buildClientCode(req);
clientCodeMap.insert(newClientCode);
log.debug("successfully processed, new clientCode id {}", newClientCode.getId());
return null;
}
protected RequestInfoUpdate clientCodeUpdate(BaseRequest<ClientCodeUpdateRequest> userRequest) {
ClientCodeUpdateRequest req = userRequest.getRequestPayload();
RequestInfoUpdate requestInfoUpdate = userRoleVerification.validateRoleAndGetResult(userRequest);
if (requestInfoUpdate != null) return requestInfoUpdate;
requestInfoUpdate = validationHelper.validateTillFirstError(userRequest, clientCodeUpdateRequestValidator);
if (requestInfoUpdate != null) return requestInfoUpdate;
log.debug("ClientCodeUpdateRequest received, id={}", req.getId());
ClientCode clientCode = clientCodeMap.getSingleObjectByID(req.getId());
if (clientCode == null) {
return makeErrorResponse(userRequest, CompanyErrors.RecordNotFound, req.getId());
}
boolean doCreateTCR = checkNeedCreateTCR(req.getCompanyId(), req.getMoneyAccountId(), req.getDepoAccountId());
if (doCreateTCR) {
try {
createAndWaitTCR(userRequest.getId(), clientCode.getId(),
req.getCompanyId(), req.getMoneyAccountId(), req.getDepoAccountId());
} catch (ClearingBaseException e) {
log.error("Can not wait creation of TCR. request id={};CompanyId={}, MoneyAccountId={}, DepoAccountId={}; {}",
userRequest.getId(), req.getCompanyId(), req.getMoneyAccountId(), req.getDepoAccountId(),
e.toString());
return makeErrorResponse(userRequest, CompanyErrors.GeneralError, "Can not create TCR at this moment");
}
}
if (req.getTradingClearingRegistryId() == null && req.getMoneyAccountId() != null) {
TradingClearingRegistry tradingClearingRegistry = selectTradingClearingRegistry(req.getCompanyId(), req.getMoneyAccountId(), req.getDepoAccountId());
if (tradingClearingRegistry == null) {
return makeErrorResponse(userRequest, CompanyErrors.TradingClearingRegistryNotFound);
} else {
req.setTradingClearingRegistryId(tradingClearingRegistry.getId());
}
}
updateClientCode(clientCode, req);
clientCodeMap.update(clientCode);
log.debug("successfully processed update, id {}", clientCode.getId());
return null;
}
private RequestInfoUpdate makeErrorResponse(BaseRequest<?> req, CompanyErrors err, Object... arg) {
String errMsg = messageResolver.resolve(new EnumMessage(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) {
log.debug("CommonDeleteRequest received id = {}", userRequest.getId());
CommonDeleteRequest req = userRequest.getRequestPayload();
RequestInfoUpdate requestInfoUpdate = userRoleVerification.validateRoleAndGetResult(userRequest);
if (requestInfoUpdate != null) return requestInfoUpdate;
requestInfoUpdate = validationHelper.validateTillFirstError(userRequest, clientCodeDeleteRequestValidator);
if (requestInfoUpdate != null) return requestInfoUpdate;
ClientCode clientCode = clientCodeMap.getSingleObjectByID(req.getId());
if (clientCode.getMoneyAccountId() != null) {
log.debug("Delete clientCode.id = {}: send message to account-service", clientCode.getId());
sendBlockTCR(clientCode.getTradingClearingRegistryId(), clientCode.getMoneyAccountId());
}
log.debug("Delete clientCode.id={}", clientCode.getId());
clientCodeMap.delete(clientCode);
return null;
}
TradingClearingRegistry selectTradingClearingRegistry(Long companyId, Long moneyAccountId, Long depoAccountId) {
Map<String, Comparable<?>> query = new HashMap<>();
query.put("companyId", companyId);
query.put("moneyAccountId", moneyAccountId);
query.put("tradingClearingRegistryType", TradingClearingRegistryType.Client_B.getKey());
if (depoAccountId != null) {
query.put("depoAaccountId", depoAccountId); // todo опечатка в поле класса, см. meta.xml!
}
TradingClearingRegistry result = tradingClearingRegistryMap.getSingleObjectByFieldValues(query);
log.trace("TradingClearingRegistry by: {}; {}found", query, result == null ? "not " : "");
return result;
}
private boolean checkNeedCreateTCR(Long companyId, Long moneyAccountId, Long depoAccountId) {
// moneyAccountId обязателен, depoAccountId опционален
if (moneyAccountId == null) {
return false;
}
TradingClearingRegistry result = selectTradingClearingRegistry(companyId, moneyAccountId, depoAccountId);
if (result == null) {
return true;
} else {
return false;
}
}
protected void createAndWaitTCR(Long reqId, Long clientCode,
Long companyId, Long moneyAccountId, Long depoAccountId) throws ClearingBaseException {
log.debug("For request {}, clientCode={} need create TCR: companyId={}, moneyAccountId={}, depoAccountId={}",
reqId, clientCode == null ? "new" : clientCode,
companyId, moneyAccountId, depoAccountId);
BaseRequest<TradingClearingRegistryNewRequest> request = new BaseRequest<>();
request.setId(idGenerator.nextId());
request.setActionType(ActionType.NEW);
TradingClearingRegistryNewRequest requestPayload = new TradingClearingRegistryNewRequest();
requestPayload.setCompanyId(companyId);
requestPayload.setMoneyAccountId(moneyAccountId);
requestPayload.setDepoAccountId(depoAccountId);
// requestPayload.setStatus(WorkflowStatus.Active.getKey());
requestPayload.setTradingClearingRegistryType(TradingClearingRegistryType.Client_B.getKey());
request.setRequestPayload(requestPayload);
try {
Long reply = accountServiceExchanger.exchange(request);
if (reply == null) {
throw new ClearingBaseException(CompanyErrors.GeneralError, "Waiting account-service timeout");
}
// примечание: ошибка и сбой (непредвиденное завершение программы) не приведёт к необратимым последствиям,
// т.к. пользователь сможет повторить запрос, а созданный на предыдущем запросе ТКР уже будет создан и найдётся.
} catch (InterruptedException e) {
throw new ClearingBaseException(CompanyErrors.GeneralError, "Waiting account-service timeout");
}
}
protected void sendBlockTCR(Long tradingClearingRegistryId, Long moneyAccountId) {
if (tradingClearingRegistryId == null) {
throw new IllegalArgumentException("tradingClearingRegistryId was null");
}
BaseRequest<TradingClearingRegistryUpdateRequest> request = new BaseRequest<>();
request.setId(idGenerator.nextId());
request.setActionType(ActionType.UPDATE);
TradingClearingRegistryUpdateRequest requestPayload = new TradingClearingRegistryUpdateRequest();
requestPayload.setId(tradingClearingRegistryId);
requestPayload.setStatus(WorkflowStatus.Blocked.getKey());
request.setRequestPayload(requestPayload);
try {
sendMessage(Consts.DESTINATION_TRADING_CLEARING_REGISTRY_UPDATE, request);
} catch (Exception e) {
log.error("Can not send block TCR message, {}", e.toString());
throw e;
}
}
/**
* @param queue Consts.*
* @param message BaseRequest
* @return
*/
private Long sendMessage(String queue, BaseRequest<?> message) {
Long sentRequestId = message.getId();
if (sentRequestId == null) {
throw new IllegalArgumentException("Message " + message.getClass().getSimpleName() + " required id.");
}
log.debug("Send to \"{}\" message {} id={}", queue, message.getActionType(), sentRequestId);
try {
Future<RecordMetadata> send = kafkaProducer.send(new ProducerRecord<>(queue, message));
send.get();
return sentRequestId;
} catch (InterruptedException e) {
log.error("Interrupt when send message, {}", ExceptionUtils.getStackTrace(e));
Thread.currentThread().interrupt();
throw new RuntimeException("Thread interrupted when send message to " + queue, e);
} catch (ExecutionException e) {
log.error("Error at send message, {} cause: {}", e.toString(), ExceptionUtils.getStackTrace(e.getCause()));
throw new RuntimeException("Send message error", e.getCause());
} catch (Exception e) {
log.error("Unexpected exception when send message: {}", ExceptionUtils.getStackTrace(e));
}
return null;
}
/**
* @param req требуется заполнить tradingClearingRegistryId по tradingClearingRegistry.code
* @return
*/
private ClientCode buildClientCode(ClientCodeNewRequest req) {
ClientCode clientCode = new ClientCode();
// clientCode.setId(idSequence.newId()); add in insert
clientCode.setCreated(Instant.now());
clientCode.setUpdated(clientCode.getCreated());
clientCode.setCompanyId(req.getCompanyId());
clientCode.setCode(req.getCode());
clientCode.setTradingClearingRegistryId(req.getTradingClearingRegistryId());
clientCode.setMoneyAccountId(req.getMoneyAccountId());
clientCode.setDepoAccountId(req.getDepoAccountId());
clientCode.setStatus(req.getStatus());
return clientCode;
}
private void updateClientCode(ClientCode clientCode, ClientCodeUpdateRequest req) {
assert clientCode.getId() != null && clientCode.getId().equals(req.getId());
clientCode.setCompanyId(req.getCompanyId());
clientCode.setCode(req.getCode());
clientCode.setTradingClearingRegistryId(req.getTradingClearingRegistryId());
clientCode.setMoneyAccountId(req.getMoneyAccountId());
clientCode.setDepoAccountId(req.getDepoAccountId());
clientCode.setStatus(req.getStatus());
clientCode.setUpdated(Instant.now());
}
}

View file

@ -0,0 +1,368 @@
package ru.spcex.clearing.company.service;
import org.apache.kafka.clients.consumer.MockConsumer;
import org.apache.kafka.clients.producer.MockProducer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.ArgumentCaptor;
import org.mockito.Captor;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.boot.test.mock.mockito.SpyBean;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit.jupiter.SpringExtension;
import ru.clearing.classes.statics.data.account.Account;
import ru.clearing.classes.statics.data.account.ClientCode;
import ru.clearing.classes.statics.data.company.Company;
import ru.clearing.classes.statics.data.profile.CompanyInfo;
import ru.clearing.classes.statics.data.registry.TradingClearingRegistry;
import ru.clearing.platform.dictionary.*;
import ru.spcex.clearing.company.config.BeanConfiguration;
import ru.spcex.clearing.company.config.validation.ClientCodeValidationConfig;
import ru.spcex.clearing.company.config.validation.ValidationConfig;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
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.ClientCodeNewRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.account.ClientCodeUpdateRequest;
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.WorkflowStatus;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import javax.annotation.PostConstruct;
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;
//todo тесты случайно ломаются из-за BiDirectionQueueExchanger. Когда там выключаю initReplyListener() то тесты стабильно работают.
@ExtendWith(SpringExtension.class)
@ContextConfiguration(classes = {
ClientCodeService.class,
ClientCodeValidationConfig.class,
ValidationConfig.class,
BeanConfiguration.class,
KafkaTestConfig.class,
ImdgTestConfig.class})
class ClientCodeServiceTest {
private static final int PARTITION = 0;
private static final Long ID = 4L;
public static final MatcherFactory.Matcher<ClientCode> CLIENT_CODE_MATCHER = usingIgnoringFieldsComparator("created","updated");
private static final Long TCR_ID = 41L;
private static final Long COMPANY_ID = 42L;
@Autowired
ClientCodeService clientCodeService;
@Autowired
@Qualifier("hazelcastServiceTest")
private ImdgProvider hazelcastServiceTest;
@Captor
private ArgumentCaptor<ProducerRecord> producerRecord;
@SpyBean
private MockProducer<String, Object> mockProducer;
private Imdg<ClientCode> clientCodeImdg;
// ****************************-*******************
@PostConstruct
private void init() {
waitAvailableImdgProviderAndAddAdminWithDefaultId();
clientCodeImdg = hazelcastServiceTest.getImdg(IMDGDistributedNames.Map_ClientCode, ClientCode.class);
// Словари для теста, применяются в ValidationConfig
putToDictionary(IMDGDistributedNames.Map_WorkflowStatusDictionary, new WorkflowStatusDictionary(), "ACTV");
putToDictionary(IMDGDistributedNames.Map_CompanySymbolDictionary, new CompanySymbolDictionary(), "CLRC");
putToDictionary(IMDGDistributedNames.Map_CorporationSoleTypeDictionary, new CorporationSoleTypeDictionary(), "GDIR");
putToDictionary(IMDGDistributedNames.Map_CountryCodeDictionary, new CountryCodeDictionary(), "RUS");
putToDictionary(IMDGDistributedNames.Map_AllowedDictionary, new AllowedDictionary(), "ALWD");
putToDictionary(IMDGDistributedNames.Map_LegalKindDictionary, new LegalKindDictionary(), "JURD");
putToDictionary(IMDGDistributedNames.Map_OrganizationTypeDictionary, new OrganizationTypeDictionary(), "NCRD");
Imdg<Company> companyImdg = hazelcastServiceTest.getImdg(IMDGDistributedNames.Map_Company, Company.class);
Company company1 = new Company();
company1.setId(COMPANY_ID);
company1.setWorkflowStatus(WorkflowStatus.Active.getKey());
company1.setFullName("Company prime");
company1.setShortName("Seizwell");
company1.setProfile(new CompanyInfo());
company1.getProfile().setCompanyId(COMPANY_ID);
company1.getProfile().setCountryCode("TLDI");
company1.getProfile().setDescription("Big profit from TLD Company Prime.");
company1.getProfile().setLegalKind("TLDI");
company1.getProfile().setResidence("TLDI");
companyImdg.insert(company1);
Imdg<TradingClearingRegistry> tcrImdg = hazelcastServiceTest.getImdg(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class);
TradingClearingRegistry registry1 = new TradingClearingRegistry();
registry1.setId(TCR_ID);
registry1.setCompanyId(COMPANY_ID);
registry1.setCode("code-120-101");
registry1.setMoneyAccountId(131L);
registry1.setDepoAaccountId(132L);
registry1.setTradingClearingRegistryType(TradingClearingRegistryType.Client_B.getKey());
registry1.setStatus(WorkflowStatus.Active.getKey());
tcrImdg.insert(registry1);
Imdg<Account> accounts = hazelcastServiceTest.getImdg(IMDGDistributedNames.Map_Account, Account.class);
Account moneyAccount=new Account();
moneyAccount.setId(131L);
moneyAccount.setAccount("AAAA-4444");
moneyAccount.setStatus("ACTV");
moneyAccount.setCompanyId(COMPANY_ID); // для валидации принадлежности счёта
accounts.insert(moneyAccount);
Account depoAccount=new Account();
depoAccount.setId(132L);
depoAccount.setAccount("AAAB-44654");
depoAccount.setStatus("ACTV");
depoAccount.setCompanyId(COMPANY_ID); // для валидации принадлежности счёта
accounts.insert(depoAccount);
TestUtils.FutureRecordMetadata future = spy(TestUtils.FutureRecordMetadata.class);
doReturn(future).when(mockProducer).send(producerRecord.capture());
}
private <D extends AbstractDictionary> void putToDictionary(String mapName, D object, String code) {
Imdg<D> dMap = (Imdg) hazelcastServiceTest.getImdg(mapName, object.getClass());
object.setId(2L);
object.setCode(code);
object.setName("name of " + code);
dMap.insert(object);
}
@Test
void selectTradingClearingRegistry() {
TradingClearingRegistry tcr = clientCodeService.selectTradingClearingRegistry(COMPANY_ID, 131L, 132L);
assertNotNull(tcr);
assertEquals(41L, tcr.getId());
tcr = clientCodeService.selectTradingClearingRegistry(COMPANY_ID, 131L, null);
assertNotNull(tcr);
assertEquals(41L, tcr.getId());
assertNull(clientCodeService.selectTradingClearingRegistry(0L, 131L, 132L));
assertNull(clientCodeService.selectTradingClearingRegistry(COMPANY_ID, 0L, 132L));
assertNull(clientCodeService.selectTradingClearingRegistry(COMPANY_ID, 131L, 0L));
}
/**
* {@link ClientCodeService#clientCodeUpdate(BaseRequest)}<br>
* Тест проверяет создание {@link ClientCode} в IMDG при передаче из Apache Kafka (очередь 1).<br>
* Входной запрос {@link ru.spcex.clearing.platform.messaging.domain.cud.account.ClientCodeNewRequest}:<br>
**/
@Test
void clientCodeNew1() {
//ARRANGE
final String ccCode = "Lucky planet";
ClientCodeNewRequest clientCodeNewRequest = new ClientCodeNewRequest();
clientCodeNewRequest.setCompanyId(COMPANY_ID);
clientCodeNewRequest.setCode(ccCode);
clientCodeNewRequest.setTradingClearingRegistryId(TCR_ID);
clientCodeNewRequest.setDepoAccountId(null);
clientCodeNewRequest.setMoneyAccountId(null);
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 = getJsonStringForNew(clientCodeNewRequest, ID);
addRecordToKafka((MockConsumer) clientCodeService.getConsumer(), Consts.DESTINATION_CLIENT_CODE_NEW, PARTITION, 0, jsonString);
//ASSERT
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 ClientCode} в IMDG при передаче из Apache Kafka (очередь 2).<br>
* Входной запрос {@link ru.spcex.clearing.platform.messaging.domain.cud.account.ClientCodeNewRequest}:<br>
**/
@Test
void clientCodeNew2() {
//ARRANGE
final String ccCode = "Lucky planet";
ClientCodeNewRequest clientCodeNewRequest = new ClientCodeNewRequest();
clientCodeNewRequest.setCompanyId(COMPANY_ID);
clientCodeNewRequest.setCode(ccCode);
clientCodeNewRequest.setTradingClearingRegistryId(TCR_ID);
clientCodeNewRequest.setDepoAccountId(null);
clientCodeNewRequest.setMoneyAccountId(null);
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 = getJsonStringForNew(clientCodeNewRequest, ID);
addRecordToKafka((MockConsumer) clientCodeService.getConsumer(), Consts.DESTINATION_CLIENT_CODE_NEW_UM_COMPANY, PARTITION, 0, jsonString);
//ASSERT
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());
}
//todo добавить тест NEW заполненными MoneyAccountId(131L), DepoAccountId(132L); - от этого направляется дополнительное сообщение в очередь и используется ожидание ответа.
/**
* {@link ClientCodeService#clientCodeUpdate(BaseRequest)}<br>
* Тест проверяет обновление сущности {@link ClientCode} в IMDG при передаче из Apache Kafka.<br>
* Входной запрос {@link ru.spcex.clearing.platform.messaging.domain.cud.account.ClientCodeUpdateRequest}:<br>
**/
@Test
void clientCodeUpdate() {
//ARRANGE
ClientCode existsClientCode = new ClientCode();
existsClientCode.setId(ID);
existsClientCode.setCompanyId(COMPANY_ID);
existsClientCode.setCode("0000");
existsClientCode.setTradingClearingRegistryId(TCR_ID);
existsClientCode.setMoneyAccountId(131L);
existsClientCode.setDepoAccountId(132L);
existsClientCode.setStatus("ACTV");
clientCodeImdg.insert(existsClientCode);
ClientCodeUpdateRequest clientCodeUpdateRequest = new ClientCodeUpdateRequest();
clientCodeUpdateRequest.setId(ID);
clientCodeUpdateRequest.setCompanyId(COMPANY_ID);
clientCodeUpdateRequest.setCode("1111");
clientCodeUpdateRequest.setTradingClearingRegistryId(TCR_ID);
clientCodeUpdateRequest.setMoneyAccountId(131L);
clientCodeUpdateRequest.setDepoAccountId(132L);
clientCodeUpdateRequest.setStatus("ACTV");
ClientCode predictableClientCode = new ClientCode();
predictableClientCode.setId(ID);
predictableClientCode.setCompanyId(COMPANY_ID);
predictableClientCode.setCode("1111");
predictableClientCode.setTradingClearingRegistryId(TCR_ID);
predictableClientCode.setMoneyAccountId(131L);
predictableClientCode.setDepoAccountId(132L);
predictableClientCode.setStatus("ACTV");
//ACT
String jsonString = getJsonStringForUpdate(clientCodeUpdateRequest, ID);
addRecordToKafka((MockConsumer) clientCodeService.getConsumer(), Consts.DESTINATION_CLIENT_CODE_UPDATE, PARTITION, 0, jsonString);
//ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord);
ClientCode resultUpdating = clientCodeImdg.getSingleObjectByID(ID);
CLIENT_CODE_MATCHER.assertMatch(resultUpdating, predictableClientCode);
assertNotNull(resultUpdating.getUpdated());
}
/**
* {@link ClientCodeService#clientCodeUpdate(BaseRequest)}<br>
* Тест проверяет удаление {@link ClientCode} из IMDG при передаче из Apache Kafka.<br>
* Входной запрос {@link ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteRequest}:<br>
**/
@Test
void clientCodeDelete1() {
//ARRANGE
ClientCode existsClientCode = new ClientCode();
existsClientCode.setId(ID);
existsClientCode.setCompanyId(COMPANY_ID);
existsClientCode.setCode("0000");
// Если следующие поля заполнить, то дополнительно отправит сообщение в trading-clearing-registry-update:
existsClientCode.setTradingClearingRegistryId(null);
existsClientCode.setMoneyAccountId(null);
existsClientCode.setDepoAccountId(null);
existsClientCode.setStatus("ACTV");
clientCodeImdg.insert(existsClientCode);
CommonDeleteRequest clientCodeDeleteRequest = new CommonDeleteRequest();
clientCodeDeleteRequest.setId(ID);
Assertions.assertNotNull(clientCodeImdg.getSingleObjectByID(ID)); // verify test data
//ACT
String jsonString = getJsonStringForUpdate(clientCodeDeleteRequest, ID);
addRecordToKafka((MockConsumer) clientCodeService.getConsumer(), Consts.DESTINATION_CLIENT_CODE_DELETE, PARTITION, 0, jsonString);
//ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord);
ClientCode resultUpdate = clientCodeImdg.getSingleObjectByID(ID);
Assertions.assertNull(resultUpdate);
}
// /**
// * {@link ClientCodeService#clientCodeUpdate(BaseRequest)}<br>
// * Тест проверяет удаление {@link ClientCode} из IMDG при передаче из Apache Kafka.<br>
// * У ClientCode заполнены MoneyAccountId, DepoAccountId - по этому при удалении должно направиться дополнительное сообщение в очередь DESTINATION_TRADING_CLEARING_REGISTRY_UPDATE<br>
// * Входной запрос {@link ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteRequest}:<br>
// **/
// @Test
// void clientCodeDelete2() {
// //ARRANGE
// ClientCode existsClientCode = new ClientCode();
// existsClientCode.setId(ID);
// existsClientCode.setCompanyId(COMPANY_ID);
// existsClientCode.setCode("0000");
// // Если следующие поля заполнить, то дополнительно отправит сообщение в trading-clearing-registry-update:
// existsClientCode.setTradingClearingRegistryId(TCR_ID);
// existsClientCode.setMoneyAccountId(131L);
// existsClientCode.setDepoAccountId(132L);
// existsClientCode.setStatus("ACTV");
//
// clientCodeImdg.insert(existsClientCode);
//
// CommonDeleteRequest clientCodeDeleteRequest = new CommonDeleteRequest();
// clientCodeDeleteRequest.setId(ID);
//
// Assertions.assertNotNull(clientCodeImdg.getSingleObjectByID(ID)); // verify test data
//
// //ACT
// String jsonString = getJsonStringForUpdate(clientCodeDeleteRequest, ID);
//
// addRecordToKafka((MockConsumer) clientCodeService.getConsumer(), Consts.DESTINATION_CLIENT_CODE_DELETE, PARTITION, 0, jsonString);
//
// //ASSERT
//
// waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord);
// ClientCode resultUpdate = clientCodeImdg.getSingleObjectByID(ID);
// Assertions.assertNull(resultUpdate);
// }
}

View file

@ -78,12 +78,14 @@ public interface Consts {
String DESTINATION_PROFILE_DOCUMENT_DELETE = "profile-document-delete";
String DESTINATION_CLIENT_CODE_NEW = "client-code-new";
String DESTINATION_CLIENT_CODE_NEW_UM_COMPANY = "client-code-new-from-api-um-company";
String DESTINATION_CLIENT_CODE_UPDATE = "client-code-update";
String DESTINATION_CLIENT_CODE_DELETE = "client-code-delete";
String DESTINATION_TRADING_CLEARING_REGISTRY_NEW = "trading-clearing-registry-new";
String DESTINATION_TRADING_CLEARING_REGISTRY_UPDATE = "trading-clearing-registry-update";
String DESTINATION_TRADING_CLEARING_REGISTRY_DELETE = "trading-clearing-registry-delete";
// String DESTINATION_TRADING_CLEARING_REGISTRY_REPLY = "trading-clearing-registry";
String DESTINATION_SDF08_NEW = "s-df-08-new";
String DESTINATION_SDF02_NEW = "s-df-02-new";

View file

@ -21,6 +21,8 @@ public class TradingClearingRegistryNewRequest {
private Long depoAccountId;
@JsonProperty
private String status;
@JsonProperty
private String tradingClearingRegistryType;
public Long getCompanyId() {
return companyId;
@ -50,6 +52,14 @@ public class TradingClearingRegistryNewRequest {
return status;
}
public String getTradingClearingRegistryType() {
return tradingClearingRegistryType;
}
public void setTradingClearingRegistryType(String tradingClearingRegistryType) {
this.tradingClearingRegistryType = tradingClearingRegistryType;
}
public void setStatus(String status) {
this.status = status;
}