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 6a0444a4e..35da485b1 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 @@ -59,11 +59,13 @@ public class BiDirectionQueueExchanger> extends Queu * @param kafkaProducer * @param outQueue отправляет в очередь * @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, + boolean ignoreOtherResponse, long timeout) { super(kafkaQueue); this.producer = kafkaProducer; @@ -85,8 +87,9 @@ public class BiDirectionQueueExchanger> extends Queu timeout = 0; } this.timeout = timeout; + this.ignoreOtherResponse = ignoreOtherResponse; - initReplyListener(); + //fixme debug only ! initReplyListener(); } /** @@ -132,6 +135,9 @@ public class BiDirectionQueueExchanger> extends Queu 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 {}", 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 index 296b3b103..0b88150b4 100644 --- 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 @@ -5,15 +5,13 @@ 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.platform.dictionary.ClearingCategoryDictionary; +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.platform.messaging.domain.cud.company.ClearingMemberCategoryNewRequest; -import ru.spcex.clearing.platform.messaging.domain.cud.company.CompanySymbolUpdateRequest; import ru.spcex.clearing.validation.common.rules.DictionaryPresentRule; import ru.spcex.clearing.validation.common.rules.FieldRequiredRule; import ru.spcex.clearing.validation.common.rules.IdPresentRule; @@ -42,7 +40,7 @@ public class ClientCodeValidationConfig { addImdg.accept(IMDGDistributedNames.Map_TradingClearingRegistry); addImdg.accept(IMDGDistributedNames.Map_WorkflowStatusDictionary); return new ValidatorImpl<>(context, - FieldRequiredRule.instance("companyId", CompanySymbolUpdateRequest::getCompanyId, CompanyErrors.RequiredFieldEmpty), + FieldRequiredRule.instance("companyId", ClientCodeNewRequest::getCompanyId, CompanyErrors.RequiredFieldEmpty), IdPresentRule.instance("companyId", ClientCodeNewRequest::getCompanyId, IMDGDistributedNames.Map_Company, @@ -73,7 +71,7 @@ public class ClientCodeValidationConfig { IdPresentRule.instance("tradingClearingRegistryId", ClientCodeNewRequest::getTradingClearingRegistryId, IMDGDistributedNames.Map_TradingClearingRegistry, - Account.class, + TradingClearingRegistry.class, CompanyErrors.RequiredFieldEmpty, CompanyErrors.TradingClearingRegistryNotFound, false, @@ -113,7 +111,7 @@ public class ClientCodeValidationConfig { CompanyErrors.RecordNotFound ), - FieldRequiredRule.instance("companyId", CompanySymbolUpdateRequest::getCompanyId, CompanyErrors.RequiredFieldEmpty), + FieldRequiredRule.instance("companyId", ClientCodeNewRequest::getCompanyId, CompanyErrors.RequiredFieldEmpty), IdPresentRule.instance("companyId", ClientCodeUpdateRequest::getCompanyId, IMDGDistributedNames.Map_Company, @@ -144,7 +142,7 @@ public class ClientCodeValidationConfig { IdPresentRule.instance("tradingClearingRegistryId", ClientCodeUpdateRequest::getTradingClearingRegistryId, IMDGDistributedNames.Map_TradingClearingRegistry, - Account.class, + TradingClearingRegistry.class, CompanyErrors.RequiredFieldEmpty, CompanyErrors.TradingClearingRegistryNotFound, false, @@ -169,7 +167,7 @@ public class ClientCodeValidationConfig { ImdgValidationContext context = new ImdgValidationContext<>(); context.setValidatedObject(companyDeleteRequest); Consumer addImdg = (s) -> context.addImdg(s, imdgForValidation.get(s)); - addImdg.accept(IMDGDistributedNames.Map_Company); + addImdg.accept(IMDGDistributedNames.Map_ClientCode); return new ValidatorImpl<>(context, // FieldRequiredRule.instance("id", CommonDeleteRequest::getId, CompanyErrors.RequiredFieldEmpty), IdPresentRule.instance("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 index a03d9c7a0..2e82f5e9b 100644 --- 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 @@ -94,7 +94,7 @@ public class ClientCodeService extends QueueConsumer implements InitializingBean this.accountServiceExchanger = new BiDirectionQueueExchanger<>(kafkaQueue, kafkaProducer, Consts.DESTINATION_TRADING_CLEARING_REGISTRY_NEW, - Consts.DESTINATION_TRADING_CLEARING_REGISTRY_REPLY, // todo naming + Consts.REQUEST_INFO_UPDATE, true, // REQUEST_INFO_UPDATE - стандартная очередь, для результатов всех реквестов. DESTINATION_TRADING_CLEARING_REGISTRY_REPLY, 60000 ); } 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 index 1510de3e8..46867f68a 100644 --- 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 @@ -39,8 +39,7 @@ import ru.spcex.platform.imdg.api.ImdgProvider; import javax.annotation.PostConstruct; import static org.junit.jupiter.api.Assertions.*; -import static org.mockito.Mockito.doReturn; -import static org.mockito.Mockito.spy; +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; @@ -179,7 +178,7 @@ class ClientCodeServiceTest { //ASSERT waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); - ClientCode resultNew = clientCodeImdg.getSingleObjectBySQL(String.format("code = %d", ccCode)); + ClientCode resultNew = clientCodeImdg.getSingleObjectBySQL(String.format("code = '%s'", ccCode)); predictableClientCode.setId(resultNew.getId()); CLIENT_CODE_MATCHER.assertMatch(resultNew, predictableClientCode); } @@ -216,7 +215,7 @@ class ClientCodeServiceTest { //ASSERT waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); - ClientCode resultNew = clientCodeImdg.getSingleObjectBySQL(String.format("code = %d", ccCode)); + ClientCode resultNew = clientCodeImdg.getSingleObjectBySQL(String.format("code = '%s'", ccCode)); predictableClientCode.setId(resultNew.getId()); CLIENT_CODE_MATCHER.assertMatch(resultNew, predictableClientCode); } @@ -276,12 +275,50 @@ class ClientCodeServiceTest { * Входной запрос {@link ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteRequest}:
**/ @Test - void clientCodeDelete() { + void clientCodeDelete1() { //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(200L); +// existsClientCode.setDepoAccountId(201L); + 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.
+ * Входной запрос {@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(200L); existsClientCode.setDepoAccountId(201L); @@ -300,10 +337,10 @@ class ClientCodeServiceTest { 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 95b708d0d..20b041d22 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 @@ -73,7 +73,7 @@ public interface Consts { 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_TRADING_CLEARING_REGISTRY_REPLY = "trading-clearing-registry"; String DESTINATION_SDF08_NEW = "s-df-08-new"; String DESTINATION_SDF02_NEW = "s-df-02-new";