diff --git a/clearing-parent/clearing-validation/src/main/java/ru/spcex/clearing/util/services/exchangers/BiDirectionQueueExchanger.java b/clearing-parent/clearing-validation/src/main/java/ru/spcex/clearing/util/services/exchangers/BiDirectionQueueExchanger.java
index f8aef7bf6..c26f63b0c 100644
--- a/clearing-parent/clearing-validation/src/main/java/ru/spcex/clearing/util/services/exchangers/BiDirectionQueueExchanger.java
+++ b/clearing-parent/clearing-validation/src/main/java/ru/spcex/clearing/util/services/exchangers/BiDirectionQueueExchanger.java
@@ -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.
+ *
+ *
+ * Ожидает из очереди (CommonIdRequest)
+ *
+ * @param отправляется в очередь
*/
-public class BiDirectionQueueExchanger> extends QueueConsumer implements InitializingBean, DisposableBean {
+public class BiDirectionQueueExchanger> 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 listenClass;
protected long timeout;
+ private Producer 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> 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 kafkaQueue, Producer kafkaProducer,
String outQueue,
- String inQueue, Class 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 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 event) {
+ synchronized (syncObject) {
+ CommonIdRequest requestPayload = event.getRequestPayload();
+ log.debug("continueWaiting {} for id={}", inQueue, requestPayload.getId());
+ lastResponseId = requestPayload.getId();
+ syncObject.notifyAll();
+ }
+ }
}
diff --git a/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/config/validation/ClientCodeValidationConfig.java b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/config/validation/ClientCodeValidationConfig.java
new file mode 100644
index 000000000..653c6320d
--- /dev/null
+++ b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/config/validation/ClientCodeValidationConfig.java
@@ -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 clientCodeNewRequestValidator(Map> imdgForValidation) {
+ return clientCodeUpdateRequest -> {
+ ImdgValidationContext context = new ImdgValidationContext<>();
+ context.setValidatedObject(clientCodeUpdateRequest);
+ Consumer 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 clientCodeUpdateRequestValidator(Map> imdgForValidation) {
+ return clientCodeUpdateRequest -> {
+ ImdgValidationContext context = new ImdgValidationContext<>();
+ context.setValidatedObject(clientCodeUpdateRequest);
+ Consumer 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>(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 clientCodeDeleteRequestValidator(Map> imdgForValidation) {
+ return companyDeleteRequest -> {
+ ImdgValidationContext context = new ImdgValidationContext<>();
+ context.setValidatedObject(companyDeleteRequest);
+ Consumer 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)
+ );
+ };
+ }
+}
diff --git a/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/config/validation/ValidationConfig.java b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/config/validation/ValidationConfig.java
index 597365270..cd6671cf1 100644
--- a/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/config/validation/ValidationConfig.java
+++ b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/config/validation/ValidationConfig.java
@@ -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;
}
diff --git a/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/error/CompanyErrors.java b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/error/CompanyErrors.java
index d36d2104f..3979ab8b9 100644
--- a/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/error/CompanyErrors.java
+++ b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/error/CompanyErrors.java
@@ -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;
diff --git a/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/service/ClientCodeService.java b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/service/ClientCodeService.java
new file mode 100644
index 000000000..2e82f5e9b
--- /dev/null
+++ b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/service/ClientCodeService.java
@@ -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 kafkaProducer;
+ private final ImdgId idGenerator;
+ private final Imdg clientCodeMap;
+ private final Imdg companyMap;
+ private final Imdg tradingClearingRegistryMap;
+
+
+ private final Function clientCodeNewRequestValidator;
+ private final Function clientCodeUpdateRequestValidator;
+ private final Function clientCodeDeleteRequestValidator;
+ private final ValidationHelper validationHelper;
+ private final UserRoleVerification userRoleVerification;
+ private final IMessageResolver messageResolver;
+
+ BiDirectionQueueExchanger> accountServiceExchanger;
+
+ @Autowired
+ public ClientCodeService(Consumer kafkaQueue,
+ Producer kafkaProducer,
+ ImdgProvider imdgProvider,
+ ValidationHelper validationHelper,
+ IMessageResolver messageResolver,
+ UserRoleVerification userRoleVerification,
+ @Qualifier("clientCodeNewRequestValidator") Function clientCodeNewRequestValidator,
+ @Qualifier("clientCodeUpdateRequestValidator") Function clientCodeUpdateRequestValidator,
+ @Qualifier("clientCodeDeleteRequestValidator") Function 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 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 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 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 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> 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 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 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 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());
+ }
+
+}
diff --git a/clearing-parent/company-service/src/test/java/ru/spcex/clearing/company/service/ClientCodeServiceTest.java b/clearing-parent/company-service/src/test/java/ru/spcex/clearing/company/service/ClientCodeServiceTest.java
new file mode 100644
index 000000000..503174d22
--- /dev/null
+++ b/clearing-parent/company-service/src/test/java/ru/spcex/clearing/company/service/ClientCodeServiceTest.java
@@ -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 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;
+ @SpyBean
+ private MockProducer mockProducer;
+
+ private Imdg 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 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 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 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 void putToDictionary(String mapName, D object, String code) {
+ Imdg 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)}
+ * Тест проверяет создание {@link ClientCode} в IMDG при передаче из Apache Kafka (очередь 1).
+ * Входной запрос {@link ru.spcex.clearing.platform.messaging.domain.cud.account.ClientCodeNewRequest}:
+ **/
+ @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)}
+ * Тест проверяет создание {@link ClientCode} в IMDG при передаче из Apache Kafka (очередь 2).
+ * Входной запрос {@link ru.spcex.clearing.platform.messaging.domain.cud.account.ClientCodeNewRequest}:
+ **/
+ @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)}
+ * Тест проверяет обновление сущности {@link ClientCode} в IMDG при передаче из Apache Kafka.
+ * Входной запрос {@link ru.spcex.clearing.platform.messaging.domain.cud.account.ClientCodeUpdateRequest}:
+ **/
+ @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)}
+ * Тест проверяет удаление {@link ClientCode} из IMDG при передаче из Apache Kafka.
+ * Входной запрос {@link ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteRequest}:
+ **/
+ @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)}
+// * Тест проверяет удаление {@link ClientCode} из IMDG при передаче из Apache Kafka.
+// * У ClientCode заполнены MoneyAccountId, DepoAccountId - по этому при удалении должно направиться дополнительное сообщение в очередь DESTINATION_TRADING_CLEARING_REGISTRY_UPDATE
+// * Входной запрос {@link ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteRequest}:
+// **/
+// @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);
+// }
+
+}
\ No newline at end of file
diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java
index c6c9f0d28..c00ea09a8 100644
--- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java
+++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java
@@ -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";
diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/registry/TradingClearingRegistryNewRequest.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/registry/TradingClearingRegistryNewRequest.java
index 75041b81e..c68eed1ba 100644
--- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/registry/TradingClearingRegistryNewRequest.java
+++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/registry/TradingClearingRegistryNewRequest.java
@@ -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;
}