Добавил поддержку нового функционала KafkaSender в тестовый модуль. Поправил тестовый модуль и тесты.

This commit is contained in:
psemenkov 2023-05-16 18:10:11 +03:00
parent c1519a2d7a
commit 4790494d62
28 changed files with 293 additions and 248 deletions

View file

@ -1,16 +1,12 @@
package ru.spcex.clearing.company.service; package ru.spcex.clearing.company.service;
import org.apache.kafka.clients.consumer.MockConsumer; import org.apache.kafka.clients.consumer.MockConsumer;
import org.apache.kafka.clients.producer.MockProducer; import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith; import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.ArgumentCaptor;
import org.mockito.Captor;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.boot.test.mock.mockito.SpyBean;
import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit.jupiter.SpringExtension; import org.springframework.test.context.junit.jupiter.SpringExtension;
import ru.clearing.classes.statics.data.company.ClearingMemberCategory; import ru.clearing.classes.statics.data.company.ClearingMemberCategory;
@ -26,7 +22,6 @@ import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteReques
import ru.spcex.clearing.platform.messaging.domain.cud.company.ClearingMemberCategoryNewRequest; import ru.spcex.clearing.platform.messaging.domain.cud.company.ClearingMemberCategoryNewRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.company.ClearingMemberCategoryUpdateRequest; import ru.spcex.clearing.platform.messaging.domain.cud.company.ClearingMemberCategoryUpdateRequest;
import ru.spcex.clearing.test.MatcherFactory; 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.ImdgTestConfig;
import ru.spcex.clearing.test.config.KafkaTestConfig; import ru.spcex.clearing.test.config.KafkaTestConfig;
import ru.spcex.platform.enumeration.ClearingCategory; import ru.spcex.platform.enumeration.ClearingCategory;
@ -36,8 +31,6 @@ import ru.spcex.platform.imdg.api.ImdgProvider;
import javax.annotation.PostConstruct; import javax.annotation.PostConstruct;
import static org.mockito.Mockito.doReturn;
import static org.mockito.Mockito.spy;
import static ru.spcex.clearing.test.MatcherFactory.usingIgnoringFieldsComparator; import static ru.spcex.clearing.test.MatcherFactory.usingIgnoringFieldsComparator;
import static ru.spcex.clearing.test.TestUtils.*; import static ru.spcex.clearing.test.TestUtils.*;
import static ru.spcex.clearing.test.config.ImdgTestConfig.currentID; import static ru.spcex.clearing.test.config.ImdgTestConfig.currentID;
@ -66,10 +59,9 @@ class ClearingMemberCategoryServiceTest {
private ImdgProvider hazelcastServiceTest; private ImdgProvider hazelcastServiceTest;
private Imdg<ClearingMemberCategory> memberCategoryImdg; private Imdg<ClearingMemberCategory> memberCategoryImdg;
@Captor @Autowired
private ArgumentCaptor<ProducerRecord> producerRecord; @Qualifier("mockProducer")
@SpyBean protected Producer<String, Object> mockProducer;
private MockProducer<String, Object> mockProducer;
private long COMPANY_ID; private long COMPANY_ID;
private Company TEST_COMPANY; private Company TEST_COMPANY;
@ -97,9 +89,6 @@ class ClearingMemberCategoryServiceTest {
clearingCategoryDictionary = new ClearingCategoryDictionary(); clearingCategoryDictionary = new ClearingCategoryDictionary();
clearingCategoryDictionary.setCode(CLEARING_CATEGORY_DICT_CODE); clearingCategoryDictionary.setCode(CLEARING_CATEGORY_DICT_CODE);
clearingCategoryDictionaryImdg.insert(clearingCategoryDictionary); clearingCategoryDictionaryImdg.insert(clearingCategoryDictionary);
TestUtils.FutureRecordMetadata future = spy(TestUtils.FutureRecordMetadata.class);
doReturn(future).when(mockProducer).send(producerRecord.capture());
} }
@Test @Test
@ -120,7 +109,7 @@ class ClearingMemberCategoryServiceTest {
addRecordToKafka((MockConsumer) clearingMemberCategoryService.getConsumer(), TOPIC_MEMBER_CATEGORY_NEW, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) clearingMemberCategoryService.getConsumer(), TOPIC_MEMBER_CATEGORY_NEW, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); waitingSendAndCheckRecord(ID, mockProducer);
ClearingMemberCategory resultNew = memberCategoryImdg.getSingleObjectBySQL(String.format("clearingMemberCategory = %s", clearingMemberCategory)); ClearingMemberCategory resultNew = memberCategoryImdg.getSingleObjectBySQL(String.format("clearingMemberCategory = %s", clearingMemberCategory));
predictableClearingMemberCategory.setId(resultNew.getId()); predictableClearingMemberCategory.setId(resultNew.getId());
MEMBER_CATEGORY_MATCHER.assertMatch(resultNew, predictableClearingMemberCategory); MEMBER_CATEGORY_MATCHER.assertMatch(resultNew, predictableClearingMemberCategory);
@ -156,7 +145,7 @@ class ClearingMemberCategoryServiceTest {
addRecordToKafka((MockConsumer) clearingMemberCategoryService.getConsumer(), TOPIC_MEMBER_CATEGORY_UPDATE, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) clearingMemberCategoryService.getConsumer(), TOPIC_MEMBER_CATEGORY_UPDATE, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(id, mockProducer, producerRecord); waitingSendAndCheckRecord(id, mockProducer);
ClearingMemberCategory resultUpdating = memberCategoryImdg.getSingleObjectByID(id); ClearingMemberCategory resultUpdating = memberCategoryImdg.getSingleObjectByID(id);
MEMBER_CATEGORY_MATCHER.assertMatch(resultUpdating, predictableClearingMemberCategory); MEMBER_CATEGORY_MATCHER.assertMatch(resultUpdating, predictableClearingMemberCategory);
@ -184,7 +173,7 @@ class ClearingMemberCategoryServiceTest {
addRecordToKafka((MockConsumer) clearingMemberCategoryService.getConsumer(), TOPIC_MEMBER_CATEGORY_DELETE, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) clearingMemberCategoryService.getConsumer(), TOPIC_MEMBER_CATEGORY_DELETE, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(id, mockProducer, producerRecord); waitingSendAndCheckRecord(id, mockProducer);
ClearingMemberCategory resultDeleting = memberCategoryImdg.getSingleObjectBySQL(String.format("clearingMemberCategory = %s", clearingMemberCategory)); ClearingMemberCategory resultDeleting = memberCategoryImdg.getSingleObjectBySQL(String.format("clearingMemberCategory = %s", clearingMemberCategory));
Assertions.assertNull(resultDeleting); Assertions.assertNull(resultDeleting);

View file

@ -1,16 +1,12 @@
package ru.spcex.clearing.company.service; package ru.spcex.clearing.company.service;
import org.apache.kafka.clients.consumer.MockConsumer; import org.apache.kafka.clients.consumer.MockConsumer;
import org.apache.kafka.clients.producer.MockProducer; import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith; import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.ArgumentCaptor;
import org.mockito.Captor;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.boot.test.mock.mockito.SpyBean;
import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit.jupiter.SpringExtension; import org.springframework.test.context.junit.jupiter.SpringExtension;
import ru.clearing.classes.statics.data.account.Account; import ru.clearing.classes.statics.data.account.Account;
@ -29,7 +25,6 @@ import ru.spcex.clearing.platform.messaging.domain.cud.account.ClientCodeNewRequ
import ru.spcex.clearing.platform.messaging.domain.cud.account.ClientCodeUpdateRequest; import ru.spcex.clearing.platform.messaging.domain.cud.account.ClientCodeUpdateRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteRequest; import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteRequest;
import ru.spcex.clearing.test.MatcherFactory; import ru.spcex.clearing.test.MatcherFactory;
import ru.spcex.clearing.test.TestUtils;
import ru.spcex.clearing.test.config.ImdgTestConfig; import ru.spcex.clearing.test.config.ImdgTestConfig;
import ru.spcex.clearing.test.config.KafkaTestConfig; import ru.spcex.clearing.test.config.KafkaTestConfig;
import ru.spcex.platform.enumeration.TradingClearingRegistryType; import ru.spcex.platform.enumeration.TradingClearingRegistryType;
@ -40,7 +35,6 @@ import ru.spcex.platform.imdg.api.ImdgProvider;
import javax.annotation.PostConstruct; import javax.annotation.PostConstruct;
import static org.junit.jupiter.api.Assertions.*; import static org.junit.jupiter.api.Assertions.*;
import static org.mockito.Mockito.*;
import static ru.spcex.clearing.test.MatcherFactory.usingIgnoringFieldsComparator; import static ru.spcex.clearing.test.MatcherFactory.usingIgnoringFieldsComparator;
import static ru.spcex.clearing.test.TestUtils.*; import static ru.spcex.clearing.test.TestUtils.*;
import static ru.spcex.clearing.test.config.ImdgTestConfig.waitAvailableImdgProviderAndAddAdminWithDefaultId; import static ru.spcex.clearing.test.config.ImdgTestConfig.waitAvailableImdgProviderAndAddAdminWithDefaultId;
@ -69,10 +63,9 @@ class ClientCodeServiceTest {
@Autowired @Autowired
@Qualifier("hazelcastServiceTest") @Qualifier("hazelcastServiceTest")
private ImdgProvider hazelcastServiceTest; private ImdgProvider hazelcastServiceTest;
@Captor @Autowired
private ArgumentCaptor<ProducerRecord> producerRecord; @Qualifier("mockProducer")
@SpyBean protected Producer<String, Object> mockProducer;
private MockProducer<String, Object> mockProducer;
private Imdg<ClientCode> clientCodeImdg; private Imdg<ClientCode> clientCodeImdg;
@ -132,9 +125,6 @@ class ClientCodeServiceTest {
depoAccount.setStatus("ACTV"); depoAccount.setStatus("ACTV");
depoAccount.setCompanyId(COMPANY_ID); // для валидации принадлежности счёта depoAccount.setCompanyId(COMPANY_ID); // для валидации принадлежности счёта
accounts.insert(depoAccount); 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) { private <D extends AbstractDictionary> void putToDictionary(String mapName, D object, String code) {
@ -192,7 +182,7 @@ class ClientCodeServiceTest {
addRecordToKafka((MockConsumer) clientCodeService.getConsumer(), Consts.DESTINATION_CLIENT_CODE_NEW, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) clientCodeService.getConsumer(), Consts.DESTINATION_CLIENT_CODE_NEW, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); waitingSendAndCheckRecord(ID, mockProducer);
ClientCode resultNew = clientCodeImdg.getSingleObjectBySQL(String.format("code = '%s'", ccCode)); ClientCode resultNew = clientCodeImdg.getSingleObjectBySQL(String.format("code = '%s'", ccCode));
predictableClientCode.setId(resultNew.getId()); predictableClientCode.setId(resultNew.getId());
CLIENT_CODE_MATCHER.assertMatch(resultNew, predictableClientCode); CLIENT_CODE_MATCHER.assertMatch(resultNew, predictableClientCode);
@ -230,7 +220,7 @@ class ClientCodeServiceTest {
addRecordToKafka((MockConsumer) clientCodeService.getConsumer(), Consts.DESTINATION_CLIENT_CODE_NEW_UM_COMPANY, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) clientCodeService.getConsumer(), Consts.DESTINATION_CLIENT_CODE_NEW_UM_COMPANY, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); waitingSendAndCheckRecord(ID, mockProducer);
ClientCode resultNew = clientCodeImdg.getSingleObjectBySQL(String.format("code = '%s'", ccCode)); ClientCode resultNew = clientCodeImdg.getSingleObjectBySQL(String.format("code = '%s'", ccCode));
predictableClientCode.setId(resultNew.getId()); predictableClientCode.setId(resultNew.getId());
CLIENT_CODE_MATCHER.assertMatch(resultNew, predictableClientCode); CLIENT_CODE_MATCHER.assertMatch(resultNew, predictableClientCode);
@ -281,7 +271,7 @@ class ClientCodeServiceTest {
addRecordToKafka((MockConsumer) clientCodeService.getConsumer(), Consts.DESTINATION_CLIENT_CODE_UPDATE, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) clientCodeService.getConsumer(), Consts.DESTINATION_CLIENT_CODE_UPDATE, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); waitingSendAndCheckRecord(ID, mockProducer);
ClientCode resultUpdating = clientCodeImdg.getSingleObjectByID(ID); ClientCode resultUpdating = clientCodeImdg.getSingleObjectByID(ID);
CLIENT_CODE_MATCHER.assertMatch(resultUpdating, predictableClientCode); CLIENT_CODE_MATCHER.assertMatch(resultUpdating, predictableClientCode);
@ -321,7 +311,7 @@ class ClientCodeServiceTest {
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); waitingSendAndCheckRecord(ID, mockProducer);
ClientCode resultUpdate = clientCodeImdg.getSingleObjectByID(ID); ClientCode resultUpdate = clientCodeImdg.getSingleObjectByID(ID);
Assertions.assertNull(resultUpdate); Assertions.assertNull(resultUpdate);
} }

View file

@ -1,15 +1,11 @@
package ru.spcex.clearing.company.service; package ru.spcex.clearing.company.service;
import org.apache.kafka.clients.consumer.MockConsumer; import org.apache.kafka.clients.consumer.MockConsumer;
import org.apache.kafka.clients.producer.MockProducer; import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith; import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.ArgumentCaptor;
import org.mockito.Captor;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.boot.test.mock.mockito.SpyBean;
import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit.jupiter.SpringExtension; import org.springframework.test.context.junit.jupiter.SpringExtension;
import ru.clearing.classes.statics.data.company.Company; import ru.clearing.classes.statics.data.company.Company;
@ -25,7 +21,6 @@ import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.platform.messaging.domain.Consts; import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.company.CompanyInfoUpdateRequest; import ru.spcex.clearing.platform.messaging.domain.cud.company.CompanyInfoUpdateRequest;
import ru.spcex.clearing.test.MatcherFactory; 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.ImdgTestConfig;
import ru.spcex.clearing.test.config.KafkaTestConfig; import ru.spcex.clearing.test.config.KafkaTestConfig;
import ru.spcex.platform.enumeration.WorkflowStatus; import ru.spcex.platform.enumeration.WorkflowStatus;
@ -34,8 +29,6 @@ import ru.spcex.platform.imdg.api.ImdgProvider;
import javax.annotation.PostConstruct; import javax.annotation.PostConstruct;
import static org.mockito.Mockito.doReturn;
import static org.mockito.Mockito.spy;
import static ru.spcex.clearing.test.MatcherFactory.usingIgnoringFieldsComparator; import static ru.spcex.clearing.test.MatcherFactory.usingIgnoringFieldsComparator;
import static ru.spcex.clearing.test.TestUtils.*; import static ru.spcex.clearing.test.TestUtils.*;
import static ru.spcex.clearing.test.config.ImdgTestConfig.waitAvailableImdgProviderAndAddAdminWithDefaultId; import static ru.spcex.clearing.test.config.ImdgTestConfig.waitAvailableImdgProviderAndAddAdminWithDefaultId;
@ -70,10 +63,9 @@ class CompanyInfoServiceTest {
private ImdgProvider hazelcastServiceTest; private ImdgProvider hazelcastServiceTest;
private Imdg<Company> companyImdg; private Imdg<Company> companyImdg;
@Captor @Autowired
private ArgumentCaptor<ProducerRecord> producerRecord; @Qualifier("mockProducer")
@SpyBean protected Producer<String, Object> mockProducer;
private MockProducer<String, Object> mockProducer;
@PostConstruct @PostConstruct
private void init() { private void init() {
@ -88,9 +80,6 @@ class CompanyInfoServiceTest {
putToDictionary(IMDGDistributedNames.Map_AllowedDictionary, new AllowedDictionary(), "ALWD"); putToDictionary(IMDGDistributedNames.Map_AllowedDictionary, new AllowedDictionary(), "ALWD");
putToDictionary(IMDGDistributedNames.Map_LegalKindDictionary, new LegalKindDictionary(), "JURD"); putToDictionary(IMDGDistributedNames.Map_LegalKindDictionary, new LegalKindDictionary(), "JURD");
putToDictionary(IMDGDistributedNames.Map_OrganizationTypeDictionary, new OrganizationTypeDictionary(), "NCRD"); putToDictionary(IMDGDistributedNames.Map_OrganizationTypeDictionary, new OrganizationTypeDictionary(), "NCRD");
TestUtils.FutureRecordMetadata future = spy(TestUtils.FutureRecordMetadata.class);
doReturn(future).when(mockProducer).send(producerRecord.capture());
} }
private <D extends AbstractDictionary /*SpcexObjectBase & Dictionary*/> void putToDictionary(String mapName, D object, String code) { private <D extends AbstractDictionary /*SpcexObjectBase & Dictionary*/> void putToDictionary(String mapName, D object, String code) {
@ -172,7 +161,7 @@ class CompanyInfoServiceTest {
addRecordToKafka((MockConsumer) companyInfoService.getConsumer(), TOPIC_COMPANY_INFO_UPDATE, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) companyInfoService.getConsumer(), TOPIC_COMPANY_INFO_UPDATE, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); waitingSendAndCheckRecord(ID, mockProducer);
Company resultUpdating = companyImdg.getSingleObjectByID(ID); Company resultUpdating = companyImdg.getSingleObjectByID(ID);
COMPANY_INFO_MATCHER.assertMatch(resultUpdating.getProfile(), predictableCompanyInfo); COMPANY_INFO_MATCHER.assertMatch(resultUpdating.getProfile(), predictableCompanyInfo);

View file

@ -1,16 +1,12 @@
package ru.spcex.clearing.company.service; package ru.spcex.clearing.company.service;
import org.apache.kafka.clients.consumer.MockConsumer; import org.apache.kafka.clients.consumer.MockConsumer;
import org.apache.kafka.clients.producer.MockProducer; import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith; import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.ArgumentCaptor;
import org.mockito.Captor;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.boot.test.mock.mockito.SpyBean;
import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit.jupiter.SpringExtension; import org.springframework.test.context.junit.jupiter.SpringExtension;
import ru.clearing.classes.statics.data.company.ClearingMemberCategory; import ru.clearing.classes.statics.data.company.ClearingMemberCategory;
@ -30,7 +26,6 @@ import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteRequest; import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.company.CompanyNewRequest; import ru.spcex.clearing.platform.messaging.domain.cud.company.CompanyNewRequest;
import ru.spcex.clearing.test.MatcherFactory; 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.ImdgTestConfig;
import ru.spcex.clearing.test.config.KafkaTestConfig; import ru.spcex.clearing.test.config.KafkaTestConfig;
import ru.spcex.platform.enumeration.CompanySymbol; import ru.spcex.platform.enumeration.CompanySymbol;
@ -42,8 +37,6 @@ import javax.annotation.PostConstruct;
import java.util.Arrays; import java.util.Arrays;
import java.util.Map; import java.util.Map;
import static org.mockito.Mockito.doReturn;
import static org.mockito.Mockito.spy;
import static ru.spcex.clearing.test.MatcherFactory.usingIgnoringFieldsComparator; import static ru.spcex.clearing.test.MatcherFactory.usingIgnoringFieldsComparator;
import static ru.spcex.clearing.test.TestUtils.*; import static ru.spcex.clearing.test.TestUtils.*;
import static ru.spcex.clearing.test.config.ImdgTestConfig.waitAvailableImdgProviderAndAddAdminWithDefaultId; import static ru.spcex.clearing.test.config.ImdgTestConfig.waitAvailableImdgProviderAndAddAdminWithDefaultId;
@ -78,10 +71,9 @@ class CompanyServiceTest {
private Imdg<CompanySymbols> companySymbolsImdg; private Imdg<CompanySymbols> companySymbolsImdg;
private Imdg<WorkflowStatusDictionary> workflowStatusDictionaryImdg; private Imdg<WorkflowStatusDictionary> workflowStatusDictionaryImdg;
@Captor @Autowired
private ArgumentCaptor<ProducerRecord> producerRecord; @Qualifier("mockProducer")
@SpyBean protected Producer<String, Object> mockProducer;
private MockProducer<String, Object> mockProducer;
@PostConstruct @PostConstruct
private void init() { private void init() {
@ -115,9 +107,6 @@ class CompanyServiceTest {
cioSymbol.setName(CompanySymbol.CIO.getKey()); cioSymbol.setName(CompanySymbol.CIO.getKey());
companySymbolDictionaryImdg.insert(cioSymbol); companySymbolDictionaryImdg.insert(cioSymbol);
} }
TestUtils.FutureRecordMetadata future = spy(TestUtils.FutureRecordMetadata.class);
doReturn(future).when(mockProducer).send(producerRecord.capture());
} }
/** /**
@ -143,7 +132,7 @@ class CompanyServiceTest {
addRecordToKafka((MockConsumer) companyService.getConsumer(), TOPIC_COMPANY_DELETE, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) companyService.getConsumer(), TOPIC_COMPANY_DELETE, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); waitingSendAndCheckRecord(ID, mockProducer);
//Company resultDeleting = companyImdg.getSingleObjectByID(ID); //Company resultDeleting = companyImdg.getSingleObjectByID(ID);
Company resultDeleting = companyImdg.getSingleObjectBySQL("id=" + ID); Company resultDeleting = companyImdg.getSingleObjectBySQL("id=" + ID);
@ -186,7 +175,7 @@ class CompanyServiceTest {
addRecordToKafka((MockConsumer) companyService.getConsumer(), DESTINATION_COMPANY_NEW, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) companyService.getConsumer(), DESTINATION_COMPANY_NEW, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); waitingSendAndCheckRecord(ID, mockProducer);
Company resultNew = companyImdg.getSingleObjectByFieldValues(Map.of("ShortName", "ClrIPO")); Company resultNew = companyImdg.getSingleObjectByFieldValues(Map.of("ShortName", "ClrIPO"));
Assertions.assertNotNull(resultNew); Assertions.assertNotNull(resultNew);
@ -221,7 +210,6 @@ class CompanyServiceTest {
//ASSERT //ASSERT
waitingWhenTryAddRecordAndCheckError(ID, waitingWhenTryAddRecordAndCheckError(ID,
mockProducer, mockProducer,
producerRecord,
String.valueOf(CompanyErrors.CompanyWithCompanySymbolAlreadyExist.getId()), String.valueOf(CompanyErrors.CompanyWithCompanySymbolAlreadyExist.getId()),
Arrays.asList("CIO", "test_value")); Arrays.asList("CIO", "test_value"));

View file

@ -1,15 +1,11 @@
package ru.spcex.clearing.company.service; package ru.spcex.clearing.company.service;
import org.apache.kafka.clients.consumer.MockConsumer; import org.apache.kafka.clients.consumer.MockConsumer;
import org.apache.kafka.clients.producer.MockProducer; import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith; import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.ArgumentCaptor;
import org.mockito.Captor;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.boot.test.mock.mockito.SpyBean;
import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit.jupiter.SpringExtension; import org.springframework.test.context.junit.jupiter.SpringExtension;
import ru.clearing.classes.statics.data.company.Company; import ru.clearing.classes.statics.data.company.Company;
@ -25,7 +21,6 @@ import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.platform.messaging.domain.Consts; import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.company.CompanySymbolUpdateRequest; import ru.spcex.clearing.platform.messaging.domain.cud.company.CompanySymbolUpdateRequest;
import ru.spcex.clearing.test.MatcherFactory; 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.ImdgTestConfig;
import ru.spcex.clearing.test.config.KafkaTestConfig; import ru.spcex.clearing.test.config.KafkaTestConfig;
import ru.spcex.platform.enumeration.CompanySymbol; import ru.spcex.platform.enumeration.CompanySymbol;
@ -34,8 +29,6 @@ import ru.spcex.platform.imdg.api.ImdgProvider;
import javax.annotation.PostConstruct; import javax.annotation.PostConstruct;
import static org.mockito.Mockito.doReturn;
import static org.mockito.Mockito.spy;
import static ru.spcex.clearing.test.MatcherFactory.usingIgnoringFieldsComparator; import static ru.spcex.clearing.test.MatcherFactory.usingIgnoringFieldsComparator;
import static ru.spcex.clearing.test.TestUtils.*; import static ru.spcex.clearing.test.TestUtils.*;
import static ru.spcex.clearing.test.config.ImdgTestConfig.waitAvailableImdgProviderAndAddAdminWithDefaultId; import static ru.spcex.clearing.test.config.ImdgTestConfig.waitAvailableImdgProviderAndAddAdminWithDefaultId;
@ -69,10 +62,9 @@ class CompanySymbolServiceTest {
private ImdgProvider hazelcastServiceTest; private ImdgProvider hazelcastServiceTest;
private Imdg<CompanySymbols> companySymbolsImdg; private Imdg<CompanySymbols> companySymbolsImdg;
@Captor @Autowired
private ArgumentCaptor<ProducerRecord> producerRecord; @Qualifier("mockProducer")
@SpyBean protected Producer<String, Object> mockProducer;
private MockProducer<String, Object> mockProducer;
@PostConstruct @PostConstruct
private void init() { private void init() {
@ -102,10 +94,6 @@ class CompanySymbolServiceTest {
company.setWorkflowStatus("ACTV"); company.setWorkflowStatus("ACTV");
companyImdg.insert(company); companyImdg.insert(company);
} }
TestUtils.FutureRecordMetadata future = spy(TestUtils.FutureRecordMetadata.class);
doReturn(future).when(mockProducer).send(producerRecord.capture());
} }
/** /**
@ -145,7 +133,7 @@ class CompanySymbolServiceTest {
addRecordToKafka((MockConsumer) companySymbolService.getConsumer(), TOPIC_COMPANY_SYMBOL_UPDATE, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) companySymbolService.getConsumer(), TOPIC_COMPANY_SYMBOL_UPDATE, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); waitingSendAndCheckRecord(ID, mockProducer);
CompanySymbols resultUpdating = companySymbolsImdg.getSingleObjectBySQL(String.format("companySymbolValue = %s", companySymbolValue)); CompanySymbols resultUpdating = companySymbolsImdg.getSingleObjectBySQL(String.format("companySymbolValue = %s", companySymbolValue));
COMPANY_SYMBOL_MATCHER.assertMatch(resultUpdating, predictableCompanySymbols); COMPANY_SYMBOL_MATCHER.assertMatch(resultUpdating, predictableCompanySymbols);

View file

@ -1,15 +1,11 @@
package ru.spcex.clearing.company.service; package ru.spcex.clearing.company.service;
import org.apache.kafka.clients.consumer.MockConsumer; import org.apache.kafka.clients.consumer.MockConsumer;
import org.apache.kafka.clients.producer.MockProducer; import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith; import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.ArgumentCaptor;
import org.mockito.Captor;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.boot.test.mock.mockito.SpyBean;
import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit.jupiter.SpringExtension; import org.springframework.test.context.junit.jupiter.SpringExtension;
import ru.clearing.classes.statics.data.company.Company; import ru.clearing.classes.statics.data.company.Company;
@ -24,7 +20,6 @@ import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.company.ContactNewRequest; import ru.spcex.clearing.platform.messaging.domain.cud.company.ContactNewRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.company.ContactUpdateRequest; import ru.spcex.clearing.platform.messaging.domain.cud.company.ContactUpdateRequest;
import ru.spcex.clearing.test.MatcherFactory; 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.ImdgTestConfig;
import ru.spcex.clearing.test.config.KafkaTestConfig; import ru.spcex.clearing.test.config.KafkaTestConfig;
import ru.spcex.platform.enumeration.ContactTypes; import ru.spcex.platform.enumeration.ContactTypes;
@ -34,8 +29,6 @@ import ru.spcex.platform.imdg.api.ImdgProvider;
import javax.annotation.PostConstruct; import javax.annotation.PostConstruct;
import static org.mockito.Mockito.doReturn;
import static org.mockito.Mockito.spy;
import static ru.spcex.clearing.test.MatcherFactory.usingIgnoringFieldsComparator; import static ru.spcex.clearing.test.MatcherFactory.usingIgnoringFieldsComparator;
import static ru.spcex.clearing.test.TestUtils.*; import static ru.spcex.clearing.test.TestUtils.*;
import static ru.spcex.clearing.test.config.ImdgTestConfig.waitAvailableImdgProviderAndAddAdminWithDefaultId; import static ru.spcex.clearing.test.config.ImdgTestConfig.waitAvailableImdgProviderAndAddAdminWithDefaultId;
@ -64,10 +57,9 @@ class ContactServiceTest {
private ImdgProvider hazelcastServiceTest; private ImdgProvider hazelcastServiceTest;
private Imdg<Contact> contactImdg; private Imdg<Contact> contactImdg;
@Captor @Autowired
private ArgumentCaptor<ProducerRecord> producerRecord; @Qualifier("mockProducer")
@SpyBean protected Producer<String, Object> mockProducer;
private MockProducer<String, Object> mockProducer;
private Company TEST_COMPANY; private Company TEST_COMPANY;
private Long COMPANY_ID; private Long COMPANY_ID;
@ -89,9 +81,6 @@ class ContactServiceTest {
IMDGDistributedNames.Map_ContactTypeDictionary, ContactTypeDictionary.class IMDGDistributedNames.Map_ContactTypeDictionary, ContactTypeDictionary.class
); );
contactTypeDictionaryImdg.insert(contactTypeDictionary); contactTypeDictionaryImdg.insert(contactTypeDictionary);
TestUtils.FutureRecordMetadata future = spy(TestUtils.FutureRecordMetadata.class);
doReturn(future).when(mockProducer).send(producerRecord.capture());
} }
@Test @Test
void contactNew() throws InterruptedException { void contactNew() throws InterruptedException {
@ -111,7 +100,7 @@ class ContactServiceTest {
addRecordToKafka((MockConsumer) contactService.getConsumer(), TOPIC_CONTACT_NEW, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) contactService.getConsumer(), TOPIC_CONTACT_NEW, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); waitingSendAndCheckRecord(ID, mockProducer);
Contact resultNew = contactImdg.getAllValues().iterator().next(); Contact resultNew = contactImdg.getAllValues().iterator().next();
predictableContact.setId(resultNew.getId()); predictableContact.setId(resultNew.getId());
@ -152,7 +141,7 @@ class ContactServiceTest {
addRecordToKafka((MockConsumer) contactService.getConsumer(), TOPIC_CONTACT_UPDATE, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) contactService.getConsumer(), TOPIC_CONTACT_UPDATE, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(existContactID, mockProducer, producerRecord); waitingSendAndCheckRecord(existContactID, mockProducer);
Contact resultUpdating = contactImdg.getSingleObjectByID(existContactID); Contact resultUpdating = contactImdg.getSingleObjectByID(existContactID);
CONTACT_MATCHER.assertMatch(resultUpdating, predictableContact); CONTACT_MATCHER.assertMatch(resultUpdating, predictableContact);

View file

@ -1,16 +1,12 @@
package ru.spcex.clearing.company.service; package ru.spcex.clearing.company.service;
import org.apache.kafka.clients.consumer.MockConsumer; import org.apache.kafka.clients.consumer.MockConsumer;
import org.apache.kafka.clients.producer.MockProducer; import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith; import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.ArgumentCaptor;
import org.mockito.Captor;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.boot.test.mock.mockito.SpyBean;
import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit.jupiter.SpringExtension; import org.springframework.test.context.junit.jupiter.SpringExtension;
import ru.clearing.classes.statics.data.company.Company; import ru.clearing.classes.statics.data.company.Company;
@ -26,7 +22,6 @@ import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteReques
import ru.spcex.clearing.platform.messaging.domain.cud.company.ProfileDocumentNewRequest; import ru.spcex.clearing.platform.messaging.domain.cud.company.ProfileDocumentNewRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.company.ProfileDocumentUpdateRequest; import ru.spcex.clearing.platform.messaging.domain.cud.company.ProfileDocumentUpdateRequest;
import ru.spcex.clearing.test.MatcherFactory; 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.ImdgTestConfig;
import ru.spcex.clearing.test.config.KafkaTestConfig; import ru.spcex.clearing.test.config.KafkaTestConfig;
import ru.spcex.platform.enumeration.DocumentTypes; import ru.spcex.platform.enumeration.DocumentTypes;
@ -38,8 +33,6 @@ import javax.annotation.PostConstruct;
import java.time.LocalDate; import java.time.LocalDate;
import java.time.temporal.ChronoUnit; import java.time.temporal.ChronoUnit;
import static org.mockito.Mockito.doReturn;
import static org.mockito.Mockito.spy;
import static ru.spcex.clearing.test.MatcherFactory.usingIgnoringFieldsComparator; import static ru.spcex.clearing.test.MatcherFactory.usingIgnoringFieldsComparator;
import static ru.spcex.clearing.test.TestUtils.*; import static ru.spcex.clearing.test.TestUtils.*;
import static ru.spcex.clearing.test.config.ImdgTestConfig.currentID; import static ru.spcex.clearing.test.config.ImdgTestConfig.currentID;
@ -71,10 +64,9 @@ class ProfileDocumentServiceTest {
private Imdg<ProfileDocument> profileDocumentMap; private Imdg<ProfileDocument> profileDocumentMap;
private Imdg<Company> companyMap; private Imdg<Company> companyMap;
@Captor @Autowired
private ArgumentCaptor<ProducerRecord> producerRecord; @Qualifier("mockProducer")
@SpyBean protected Producer<String, Object> mockProducer;
private MockProducer<String, Object> mockProducer;
private long COMPANY_ID; private long COMPANY_ID;
private Company TEST_COMPANY; private Company TEST_COMPANY;
@ -111,9 +103,6 @@ class ProfileDocumentServiceTest {
documentTypeDictionary1.setCode(NEW_TEST_DOCUMENT_TYPE.getKey()); documentTypeDictionary1.setCode(NEW_TEST_DOCUMENT_TYPE.getKey());
documentTypeDictionary1.setName("new_test_name"); documentTypeDictionary1.setName("new_test_name");
documentTypeDictionaryImdg.insert(documentTypeDictionary1); documentTypeDictionaryImdg.insert(documentTypeDictionary1);
TestUtils.FutureRecordMetadata future = spy(TestUtils.FutureRecordMetadata.class);
doReturn(future).when(mockProducer).send(producerRecord.capture());
} }
/** /**
@ -159,7 +148,7 @@ class ProfileDocumentServiceTest {
addRecordToKafka((MockConsumer) profileDocumentService.getConsumer(), TOPIC_DESTINATION_PROFILE_DOCUMENT_NEW, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) profileDocumentService.getConsumer(), TOPIC_DESTINATION_PROFILE_DOCUMENT_NEW, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); waitingSendAndCheckRecord(ID, mockProducer);
ProfileDocument resultNew = profileDocumentMap.getSingleObjectBySQL(String.format("companyId = %d", COMPANY_ID)); ProfileDocument resultNew = profileDocumentMap.getSingleObjectBySQL(String.format("companyId = %d", COMPANY_ID));
predictableProfileDocument.setId(resultNew.getId()); predictableProfileDocument.setId(resultNew.getId());
PROFILE_DOCUMENT_MATCHER.assertMatch(resultNew, predictableProfileDocument); PROFILE_DOCUMENT_MATCHER.assertMatch(resultNew, predictableProfileDocument);
@ -220,7 +209,7 @@ class ProfileDocumentServiceTest {
addRecordToKafka((MockConsumer) profileDocumentService.getConsumer(), TOPIC_DESTINATION_PROFILE_DOCUMENT_UPDATE, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) profileDocumentService.getConsumer(), TOPIC_DESTINATION_PROFILE_DOCUMENT_UPDATE, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); waitingSendAndCheckRecord(ID, mockProducer);
ProfileDocument resultUpdate = profileDocumentMap.getSingleObjectBySQL(String.format("companyId = %d", COMPANY_ID)); ProfileDocument resultUpdate = profileDocumentMap.getSingleObjectBySQL(String.format("companyId = %d", COMPANY_ID));
// predictableProfileDocument.setId(resultUpdate.getId()); // predictableProfileDocument.setId(resultUpdate.getId());
PROFILE_DOCUMENT_MATCHER.assertMatch(resultUpdate, predictableProfileDocument); PROFILE_DOCUMENT_MATCHER.assertMatch(resultUpdate, predictableProfileDocument);
@ -255,7 +244,7 @@ class ProfileDocumentServiceTest {
addRecordToKafka((MockConsumer) profileDocumentService.getConsumer(), TOPIC_DESTINATION_PROFILE_DOCUMENT_DELETE, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) profileDocumentService.getConsumer(), TOPIC_DESTINATION_PROFILE_DOCUMENT_DELETE, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); waitingSendAndCheckRecord(ID, mockProducer);
ProfileDocument resultUpdate = profileDocumentMap.getSingleObjectByID(PROFILE_DOCUMENT_ID); ProfileDocument resultUpdate = profileDocumentMap.getSingleObjectByID(PROFILE_DOCUMENT_ID);
Assertions.assertNull(resultUpdate); Assertions.assertNull(resultUpdate);
} }

View file

@ -1,15 +1,11 @@
package ru.spcex.clearing.company.service; package ru.spcex.clearing.company.service;
import org.apache.kafka.clients.consumer.MockConsumer; import org.apache.kafka.clients.consumer.MockConsumer;
import org.apache.kafka.clients.producer.MockProducer; import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith; import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.ArgumentCaptor;
import org.mockito.Captor;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.boot.test.mock.mockito.SpyBean;
import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit.jupiter.SpringExtension; import org.springframework.test.context.junit.jupiter.SpringExtension;
import ru.clearing.classes.statics.data.company.Company; import ru.clearing.classes.statics.data.company.Company;
@ -25,7 +21,6 @@ import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteReques
import ru.spcex.clearing.platform.messaging.domain.cud.company.RelationNewRequest; import ru.spcex.clearing.platform.messaging.domain.cud.company.RelationNewRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.company.RelationUpdateRequest; import ru.spcex.clearing.platform.messaging.domain.cud.company.RelationUpdateRequest;
import ru.spcex.clearing.test.MatcherFactory; 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.ImdgTestConfig;
import ru.spcex.clearing.test.config.KafkaTestConfig; import ru.spcex.clearing.test.config.KafkaTestConfig;
import ru.spcex.platform.enumeration.ClearingCategory; import ru.spcex.platform.enumeration.ClearingCategory;
@ -35,9 +30,7 @@ import ru.spcex.platform.imdg.api.ImdgProvider;
import javax.annotation.PostConstruct; import javax.annotation.PostConstruct;
import static org.junit.jupiter.api.Assertions.*; import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.mockito.Mockito.doReturn;
import static org.mockito.Mockito.spy;
import static ru.spcex.clearing.test.MatcherFactory.usingIgnoringFieldsComparator; import static ru.spcex.clearing.test.MatcherFactory.usingIgnoringFieldsComparator;
import static ru.spcex.clearing.test.TestUtils.*; import static ru.spcex.clearing.test.TestUtils.*;
import static ru.spcex.clearing.test.config.ImdgTestConfig.currentID; import static ru.spcex.clearing.test.config.ImdgTestConfig.currentID;
@ -66,10 +59,9 @@ class RelationServiceTest {
private Imdg<Relation> relationMap; private Imdg<Relation> relationMap;
private Imdg<Company> companyMap; private Imdg<Company> companyMap;
@Captor @Autowired
private ArgumentCaptor<ProducerRecord> producerRecord; @Qualifier("mockProducer")
@SpyBean protected Producer<String, Object> mockProducer;
private MockProducer<String, Object> mockProducer;
private long COMPANY_ID; private long COMPANY_ID;
@ -109,15 +101,13 @@ class RelationServiceTest {
clearingCategoryDictionary.setCode(ClearingCategory.I.getKey()); clearingCategoryDictionary.setCode(ClearingCategory.I.getKey());
clearingCategoryDictionary.setName("I"); clearingCategoryDictionary.setName("I");
clearingCategoryDictionaryImdg.insert(clearingCategoryDictionary); clearingCategoryDictionaryImdg.insert(clearingCategoryDictionary);
TestUtils.FutureRecordMetadata future = spy(TestUtils.FutureRecordMetadata.class);
doReturn(future).when(mockProducer).send(producerRecord.capture());
} }
@Test @Test
void newRelation() { void newRelation() {
//ARRANGE //ARRANGE
clearAllInImdg(relationMap);
RelationNewRequest relationNewRequest = new RelationNewRequest(); RelationNewRequest relationNewRequest = new RelationNewRequest();
relationNewRequest.setCompanyId(COMPANY_ID); relationNewRequest.setCompanyId(COMPANY_ID);
relationNewRequest.setServiceStatus("ACTV"); relationNewRequest.setServiceStatus("ACTV");
@ -138,7 +128,7 @@ class RelationServiceTest {
addRecordToKafka((MockConsumer) relationService.getConsumer(), Consts.DESTINATION_RELATION_NEW, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) relationService.getConsumer(), Consts.DESTINATION_RELATION_NEW, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); waitingSendAndCheckRecord(ID, mockProducer);
Relation resultNew = relationMap.getSingleObjectBySQL("serviceStatus=ACTV"); // String.format("supplierId = %d", COMPANY_ID)); Relation resultNew = relationMap.getSingleObjectBySQL("serviceStatus=ACTV"); // String.format("supplierId = %d", COMPANY_ID));
predictableRelation.setId(resultNew.getId()); predictableRelation.setId(resultNew.getId());
RELATION_MATCHER.assertMatch(resultNew, predictableRelation); RELATION_MATCHER.assertMatch(resultNew, predictableRelation);
@ -177,7 +167,7 @@ class RelationServiceTest {
addRecordToKafka((MockConsumer) relationService.getConsumer(), Consts.DESTINATION_RELATION_UPDATE, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) relationService.getConsumer(), Consts.DESTINATION_RELATION_UPDATE, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); waitingSendAndCheckRecord(ID, mockProducer);
Relation resultUpdate = relationMap.getSingleObjectBySQL(String.format("id = %d", RELATION_ID)); Relation resultUpdate = relationMap.getSingleObjectBySQL(String.format("id = %d", RELATION_ID));
// predictableProfileDocument.setId(resultUpdate.getId()); // predictableProfileDocument.setId(resultUpdate.getId());
RELATION_MATCHER.assertMatch(resultUpdate, predictableRelation); RELATION_MATCHER.assertMatch(resultUpdate, predictableRelation);
@ -206,7 +196,7 @@ class RelationServiceTest {
addRecordToKafka((MockConsumer) relationService.getConsumer(), Consts.DESTINATION_RELATION_DELETE, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) relationService.getConsumer(), Consts.DESTINATION_RELATION_DELETE, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); waitingSendAndCheckRecord(ID, mockProducer);
Relation resultUpdate = relationMap.getSingleObjectByID(RELATION_ID); Relation resultUpdate = relationMap.getSingleObjectByID(RELATION_ID);
// Assertions.assertNull(resultUpdate); // Assertions.assertNull(resultUpdate);
assertEquals(WorkflowStatus.Blocked.getKey(), resultUpdate.getServiceStatus()); assertEquals(WorkflowStatus.Blocked.getKey(), resultUpdate.getServiceStatus());

View file

@ -1,13 +1,9 @@
package ru.spcex.clearing.scheduler; package ru.spcex.clearing.scheduler;
import org.apache.kafka.clients.producer.MockProducer; import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.junit.jupiter.api.extension.ExtendWith; 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.Autowired;
import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.boot.test.mock.mockito.MockBean;
import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit.jupiter.SpringExtension; import org.springframework.test.context.junit.jupiter.SpringExtension;
import ru.clearing.classes.statics.data.company.Company; import ru.clearing.classes.statics.data.company.Company;
@ -27,7 +23,6 @@ import ru.spcex.clearing.scheduler.config.validation.PlannerValidationConfig;
import ru.spcex.clearing.scheduler.config.validation.ValidationConfig; import ru.spcex.clearing.scheduler.config.validation.ValidationConfig;
import ru.spcex.clearing.scheduler.service.*; import ru.spcex.clearing.scheduler.service.*;
import ru.spcex.clearing.test.MatcherFactory; 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.ImdgTestConfig;
import ru.spcex.clearing.test.config.KafkaTestConfig; import ru.spcex.clearing.test.config.KafkaTestConfig;
import ru.spcex.platform.enumeration.*; import ru.spcex.platform.enumeration.*;
@ -37,8 +32,6 @@ import ru.spcex.platform.imdg.api.ImdgProvider;
import java.time.LocalDate; import java.time.LocalDate;
import java.time.LocalTime; import java.time.LocalTime;
import static org.mockito.Mockito.doReturn;
import static org.mockito.Mockito.spy;
import static ru.spcex.clearing.test.MatcherFactory.usingIgnoringFieldsComparator; import static ru.spcex.clearing.test.MatcherFactory.usingIgnoringFieldsComparator;
import static ru.spcex.clearing.test.config.ImdgTestConfig.currentID; import static ru.spcex.clearing.test.config.ImdgTestConfig.currentID;
import static ru.spcex.clearing.test.config.ImdgTestConfig.waitAvailableImdgProviderAndAddAdminWithDefaultId; import static ru.spcex.clearing.test.config.ImdgTestConfig.waitAvailableImdgProviderAndAddAdminWithDefaultId;
@ -74,10 +67,10 @@ public abstract class AbstractServiceTest {
protected long testSecurityId = currentID.getAndIncrement(); protected long testSecurityId = currentID.getAndIncrement();
protected String sessionType = "IPOT"; protected String sessionType = "IPOT";
@Captor @Autowired
protected ArgumentCaptor<ProducerRecord> producerRecord; @Qualifier("mockProducer")
@MockBean protected Producer<String, Object> mockProducer;
protected MockProducer<String, Object> mockProducer;
@Autowired @Autowired
@Qualifier("hazelcastServiceTest") @Qualifier("hazelcastServiceTest")
protected ImdgProvider imdgProvider; protected ImdgProvider imdgProvider;
@ -133,9 +126,6 @@ public abstract class AbstractServiceTest {
securityImdg.insert(security); securityImdg.insert(security);
clearingCalendarImdg.insert(getClearingCalendar()); clearingCalendarImdg.insert(getClearingCalendar());
TestUtils.FutureRecordMetadata future = spy(new TestUtils.FutureRecordMetadata());
doReturn(future).when(mockProducer).send(producerRecord.capture());
} }
protected PlannerTemplate getPlannerTemplate(long companyId) { protected PlannerTemplate getPlannerTemplate(long companyId) {

View file

@ -75,7 +75,7 @@ class ClearingCalendarServiceTest extends AbstractServiceTest {
addRecordToKafka((MockConsumer) clearingCalendarService.getConsumer(), TOPIC_CLEARING_CALENDAR_NEW, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) clearingCalendarService.getConsumer(), TOPIC_CLEARING_CALENDAR_NEW, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); waitingSendAndCheckRecord(ID, mockProducer);
ClearingCalendar clearingCalendarRes = clearingCalendarImdg.getSingleObjectBySQL(String.format("companyId = %s", clearingCalendar.getCompanyId())); ClearingCalendar clearingCalendarRes = clearingCalendarImdg.getSingleObjectBySQL(String.format("companyId = %s", clearingCalendar.getCompanyId()));
clearingCalendar.setId(clearingCalendarRes.getId()); clearingCalendar.setId(clearingCalendarRes.getId());
@ -113,7 +113,7 @@ class ClearingCalendarServiceTest extends AbstractServiceTest {
addRecordToKafka((MockConsumer) clearingCalendarService.getConsumer(), TOPIC_CLEARING_CALENDAR_UPDATE, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) clearingCalendarService.getConsumer(), TOPIC_CLEARING_CALENDAR_UPDATE, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); waitingSendAndCheckRecord(ID, mockProducer);
ClearingCalendar clearingCalendarRes = clearingCalendarImdg.getSingleObjectBySQL(String.format("companyId = %s", clearingCalendar.getCompanyId())); ClearingCalendar clearingCalendarRes = clearingCalendarImdg.getSingleObjectBySQL(String.format("companyId = %s", clearingCalendar.getCompanyId()));
clearingCalendar.setId(clearingCalendarRes.getId()); clearingCalendar.setId(clearingCalendarRes.getId());
@ -148,7 +148,7 @@ class ClearingCalendarServiceTest extends AbstractServiceTest {
addRecordToKafka((MockConsumer) clearingCalendarService.getConsumer(), TOPIC_CLEARING_CALENDAR_DELETE, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) clearingCalendarService.getConsumer(), TOPIC_CLEARING_CALENDAR_DELETE, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(id, mockProducer, producerRecord); waitingSendAndCheckRecord(id, mockProducer);
ClearingCalendar plannerTemplateRes = clearingCalendarImdg.getSingleObjectBySQL(String.format("companyId = %s", clearingCalendar.getCompanyId())); ClearingCalendar plannerTemplateRes = clearingCalendarImdg.getSingleObjectBySQL(String.format("companyId = %s", clearingCalendar.getCompanyId()));
assertNull(plannerTemplateRes); assertNull(plannerTemplateRes);

View file

@ -1,8 +1,10 @@
package ru.spcex.clearing.scheduler.service; package ru.spcex.clearing.scheduler.service;
import org.apache.kafka.clients.consumer.MockConsumer; import org.apache.kafka.clients.consumer.MockConsumer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Test;
import org.mockito.ArgumentCaptor;
import org.slf4j.Logger; import org.slf4j.Logger;
import org.slf4j.LoggerFactory; import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
@ -24,6 +26,7 @@ import static ru.spcex.clearing.test.MatcherFactory.usingIgnoringFieldsComparato
import static ru.spcex.clearing.test.TestUtils.*; import static ru.spcex.clearing.test.TestUtils.*;
import static ru.spcex.clearing.test.config.ImdgTestConfig.currentID; import static ru.spcex.clearing.test.config.ImdgTestConfig.currentID;
import static ru.spcex.clearing.test.config.ImdgTestConfig.defaultAdminId; import static ru.spcex.clearing.test.config.ImdgTestConfig.defaultAdminId;
import static ru.spcex.clearing.test.config.KafkaTestConfig.getCaptor;
import static ru.spcex.platform.enumeration.Task.accountBlock; import static ru.spcex.platform.enumeration.Task.accountBlock;
class LauncherServiceTest extends AbstractServiceTest { class LauncherServiceTest extends AbstractServiceTest {
@ -76,11 +79,12 @@ class LauncherServiceTest extends AbstractServiceTest {
addRecordToKafka((MockConsumer) launcherService.getConsumer(), TOPIC_LAUNCHER_NEW, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) launcherService.getConsumer(), TOPIC_LAUNCHER_NEW, PARTITION, 0, jsonString);
//ASSERT //ASSERT
ArgumentCaptor<ProducerRecord> captor = getCaptor(mockProducer);
verify(mockProducer, timeout(30_000L).times(1)) verify(mockProducer, timeout(30_000L).times(1))
.send(producerRecord.capture()); .send(captor.capture());
BaseRequest<Object> baseRequestResult = (BaseRequest<Object>) producerRecord.getValue().value(); BaseRequest<Object> baseRequestResult = (BaseRequest<Object>) captor.getValue().value();
assertEquals("launcher-" + launcher.getTask(), producerRecord.getValue().topic()); assertEquals("launcher-" + launcher.getTask(), captor.getValue().topic());
BASE_REQUEST_MATCHER.assertMatch(baseRequestResult, predictableBaseRequest); BASE_REQUEST_MATCHER.assertMatch(baseRequestResult, predictableBaseRequest);

View file

@ -69,7 +69,7 @@ class PlannerServiceTest extends AbstractServiceTest {
addRecordToKafka((MockConsumer) plannerService.getConsumer(), TOPIC_PLANNER_NEW, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) plannerService.getConsumer(), TOPIC_PLANNER_NEW, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); waitingSendAndCheckRecord(ID, mockProducer);
Planner plannerReq = plannerImdg.getSingleObjectBySQL(String.format("companyId = %s", planner.getCompanyId())); Planner plannerReq = plannerImdg.getSingleObjectBySQL(String.format("companyId = %s", planner.getCompanyId()));
planner.setId(plannerReq.getId()); planner.setId(plannerReq.getId());
@ -105,7 +105,7 @@ class PlannerServiceTest extends AbstractServiceTest {
addRecordToKafka((MockConsumer) plannerService.getConsumer(), TOPIC_PLANNER_UPDATE, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) plannerService.getConsumer(), TOPIC_PLANNER_UPDATE, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(planner.getId(), mockProducer, producerRecord); waitingSendAndCheckRecord(planner.getId(), mockProducer);
Planner plannerReq = plannerImdg.getSingleObjectBySQL(String.format("companyId = %s", planner.getCompanyId())); Planner plannerReq = plannerImdg.getSingleObjectBySQL(String.format("companyId = %s", planner.getCompanyId()));
planner.setId(plannerReq.getId()); planner.setId(plannerReq.getId());
@ -137,7 +137,7 @@ class PlannerServiceTest extends AbstractServiceTest {
addRecordToKafka((MockConsumer) plannerService.getConsumer(), TOPIC_PLANNER_DELETE, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) plannerService.getConsumer(), TOPIC_PLANNER_DELETE, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(plannerId, mockProducer, producerRecord); waitingSendAndCheckRecord(plannerId, mockProducer);
Planner plannerReq = plannerImdg.getSingleObjectBySQL(String.format("companyId = %s", planner.getCompanyId())); Planner plannerReq = plannerImdg.getSingleObjectBySQL(String.format("companyId = %s", planner.getCompanyId()));
assertNull(plannerReq); assertNull(plannerReq);

View file

@ -61,7 +61,7 @@ class PlannerTemplateServiceTest extends AbstractServiceTest {
addRecordToKafka((MockConsumer) plannerTemplateService.getConsumer(), TOPIC_PLANNER_TEMPLATE_NEW, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) plannerTemplateService.getConsumer(), TOPIC_PLANNER_TEMPLATE_NEW, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); waitingSendAndCheckRecord(ID, mockProducer);
PlannerTemplate plannerTemplateRes = plannerTemplateImdg.getSingleObjectBySQL(String.format("companyId = %s", plannerTemplate.getCompanyId())); PlannerTemplate plannerTemplateRes = plannerTemplateImdg.getSingleObjectBySQL(String.format("companyId = %s", plannerTemplate.getCompanyId()));
plannerTemplate.setId(plannerTemplateRes.getId()); plannerTemplate.setId(plannerTemplateRes.getId());
@ -96,7 +96,7 @@ class PlannerTemplateServiceTest extends AbstractServiceTest {
addRecordToKafka((MockConsumer) plannerTemplateService.getConsumer(), TOPIC_PLANNER_TEMPLATE_UPDATE, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) plannerTemplateService.getConsumer(), TOPIC_PLANNER_TEMPLATE_UPDATE, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); waitingSendAndCheckRecord(ID, mockProducer);
PlannerTemplate plannerTemplateRes = plannerTemplateImdg.getSingleObjectBySQL(String.format("companyId = %s", plannerTemplate.getCompanyId())); PlannerTemplate plannerTemplateRes = plannerTemplateImdg.getSingleObjectBySQL(String.format("companyId = %s", plannerTemplate.getCompanyId()));
plannerTemplate.setId(plannerTemplateRes.getId()); plannerTemplate.setId(plannerTemplateRes.getId());
@ -127,7 +127,7 @@ class PlannerTemplateServiceTest extends AbstractServiceTest {
addRecordToKafka((MockConsumer) plannerTemplateService.getConsumer(), TOPIC_PLANNER_TEMPLATE_DELETE, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) plannerTemplateService.getConsumer(), TOPIC_PLANNER_TEMPLATE_DELETE, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(id, mockProducer, producerRecord); waitingSendAndCheckRecord(id, mockProducer);
PlannerTemplate plannerTemplateRes = plannerTemplateImdg.getSingleObjectBySQL(String.format("companyId = %s", plannerTemplate.getCompanyId())); PlannerTemplate plannerTemplateRes = plannerTemplateImdg.getSingleObjectBySQL(String.format("companyId = %s", plannerTemplate.getCompanyId()));
assertNull(plannerTemplateRes); assertNull(plannerTemplateRes);

View file

@ -1,6 +1,8 @@
package ru.spcex.clearing.scheduler.service; package ru.spcex.clearing.scheduler.service;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Test;
import org.mockito.ArgumentCaptor;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import ru.clearing.classes.statics.data.scheduler.PlannerAllToday; import ru.clearing.classes.statics.data.scheduler.PlannerAllToday;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest; import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
@ -18,6 +20,7 @@ import static ru.spcex.clearing.scheduler.config.PlannerQueueConfig.addToPlanner
import static ru.spcex.clearing.scheduler.service.TaskManager.systemId; import static ru.spcex.clearing.scheduler.service.TaskManager.systemId;
import static ru.spcex.clearing.test.TestUtils.BASE_REQUEST_MATCHER; import static ru.spcex.clearing.test.TestUtils.BASE_REQUEST_MATCHER;
import static ru.spcex.clearing.test.config.ImdgTestConfig.currentID; import static ru.spcex.clearing.test.config.ImdgTestConfig.currentID;
import static ru.spcex.clearing.test.config.KafkaTestConfig.getCaptor;
import static ru.spcex.platform.utils.enumeration.IEnumKey.getEnumByKey; import static ru.spcex.platform.utils.enumeration.IEnumKey.getEnumByKey;
class TaskManagerTest extends AbstractServiceTest { class TaskManagerTest extends AbstractServiceTest {
@ -25,6 +28,9 @@ class TaskManagerTest extends AbstractServiceTest {
@Autowired @Autowired
LauncherSender launcherSender; LauncherSender launcherSender;
@Autowired
private LauncherService launcherService;
@PostConstruct @PostConstruct
public void init() { public void init() {
super.init(); super.init();
@ -95,8 +101,9 @@ class TaskManagerTest extends AbstractServiceTest {
public void waitingWhenAddedLauncherCommandRequestAndCheckIt(Task toTaskQueue, Long userId) { public void waitingWhenAddedLauncherCommandRequestAndCheckIt(Task toTaskQueue, Long userId) {
BaseRequest<Object> predictableBaseRequest = launcherSender.makeCmdRequest(toTaskQueue, userId); BaseRequest<Object> predictableBaseRequest = launcherSender.makeCmdRequest(toTaskQueue, userId);
ArgumentCaptor<ProducerRecord> producerRecord = getCaptor(mockProducer);
//waiting for kafka producer send message (finale event) //waiting for kafka producer send message (finale event)
verify(mockProducer, timeout(60_000L).times(1)) verify(mockProducer, timeout(30_000L).times(1))
.send(producerRecord.capture()); .send(producerRecord.capture());
BaseRequest<Object> baseRequestResult = (BaseRequest<Object>) producerRecord.getValue().value(); BaseRequest<Object> baseRequestResult = (BaseRequest<Object>) producerRecord.getValue().value();

View file

@ -1,13 +1,9 @@
package ru.spcex.clearing.securities.service; package ru.spcex.clearing.securities.service;
import org.apache.kafka.clients.producer.MockProducer; import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.junit.jupiter.api.extension.ExtendWith; 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.Autowired;
import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.boot.test.mock.mockito.MockBean;
import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit.jupiter.SpringExtension; import org.springframework.test.context.junit.jupiter.SpringExtension;
import ru.clearing.classes.statics.data.company.Company; import ru.clearing.classes.statics.data.company.Company;
@ -24,7 +20,6 @@ import ru.spcex.clearing.securities.config.ValidationConfig;
import ru.spcex.clearing.securities.service.cud.*; import ru.spcex.clearing.securities.service.cud.*;
import ru.spcex.clearing.securities.validation.ValidationProvider; import ru.spcex.clearing.securities.validation.ValidationProvider;
import ru.spcex.clearing.test.MatcherFactory; 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.ImdgTestConfig;
import ru.spcex.clearing.test.config.KafkaTestConfig; import ru.spcex.clearing.test.config.KafkaTestConfig;
import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.Imdg;
@ -32,8 +27,6 @@ import ru.spcex.platform.imdg.api.ImdgProvider;
import java.util.concurrent.atomic.AtomicLong; import java.util.concurrent.atomic.AtomicLong;
import static org.mockito.Mockito.doReturn;
import static org.mockito.Mockito.spy;
import static ru.spcex.clearing.test.MatcherFactory.usingIgnoringFieldsComparator; import static ru.spcex.clearing.test.MatcherFactory.usingIgnoringFieldsComparator;
import static ru.spcex.clearing.test.config.ImdgTestConfig.waitAvailableImdgProviderAndAddAdminWithDefaultId; import static ru.spcex.clearing.test.config.ImdgTestConfig.waitAvailableImdgProviderAndAddAdminWithDefaultId;
@ -75,10 +68,9 @@ public abstract class AbstractServiceTest {
public static String termType = "TermType"; public static String termType = "TermType";
public static String currencyCode = "RUB"; public static String currencyCode = "RUB";
@Captor @Autowired
protected ArgumentCaptor<ProducerRecord> producerRecord; @Qualifier("mockProducer")
@MockBean protected Producer<String, Object> mockProducer;
protected MockProducer<String, Object> mockProducer;
@Autowired @Autowired
@Qualifier("hazelcastServiceTest") @Qualifier("hazelcastServiceTest")
protected ImdgProvider imdgProvider; protected ImdgProvider imdgProvider;
@ -127,7 +119,5 @@ public abstract class AbstractServiceTest {
Company company = new Company(); Company company = new Company();
company.setId(issuerId); company.setId(issuerId);
companyImdg.insert(company); companyImdg.insert(company);
TestUtils.FutureRecordMetadata future = spy(TestUtils.FutureRecordMetadata.class);
doReturn(future).when(mockProducer).send(producerRecord.capture());
} }
} }

View file

@ -63,7 +63,7 @@ public class CouponPeriodServiceTest extends AbstractServiceTest {
addRecordToKafka((MockConsumer) couponPeriodService.getConsumer(), TOPIC, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) couponPeriodService.getConsumer(), TOPIC, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); waitingSendAndCheckRecord(ID, mockProducer);
CouponPeriod result = couponPeriodImdg.getSingleObjectBySQL(String.format("securityId = %s", currencyPrediction.getSecurityId())); CouponPeriod result = couponPeriodImdg.getSingleObjectBySQL(String.format("securityId = %s", currencyPrediction.getSecurityId()));
currencyPrediction.setId(result.getId()); currencyPrediction.setId(result.getId());
@ -112,7 +112,7 @@ public class CouponPeriodServiceTest extends AbstractServiceTest {
addRecordToKafka((MockConsumer) couponPeriodService.getConsumer(), TOPIC, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) couponPeriodService.getConsumer(), TOPIC, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(id, mockProducer, producerRecord); waitingSendAndCheckRecord(id, mockProducer);
CouponPeriod result = couponPeriodImdg.getSingleObjectBySQL(String.format("securityId = %s", predictableCoupon.getSecurityId())); CouponPeriod result = couponPeriodImdg.getSingleObjectBySQL(String.format("securityId = %s", predictableCoupon.getSecurityId()));
predictableCoupon.setId(result.getId()); predictableCoupon.setId(result.getId());

View file

@ -54,7 +54,7 @@ public class CurrencyServiceTest extends AbstractServiceTest {
addRecordToKafka((MockConsumer) currencyService.getConsumer(), TOPIC, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) currencyService.getConsumer(), TOPIC, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); waitingSendAndCheckRecord(ID, mockProducer);
Currency result = currencyImdg.getSingleObjectBySQL(String.format("countryCode = %s", currencyPrediction.getCountryCode())); Currency result = currencyImdg.getSingleObjectBySQL(String.format("countryCode = %s", currencyPrediction.getCountryCode()));
currencyPrediction.setId(result.getId()); currencyPrediction.setId(result.getId());
@ -87,7 +87,7 @@ public class CurrencyServiceTest extends AbstractServiceTest {
addRecordToKafka((MockConsumer) currencyService.getConsumer(), TOPIC, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) currencyService.getConsumer(), TOPIC, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); waitingSendAndCheckRecord(ID, mockProducer);
Currency result = currencyImdg.getSingleObjectBySQL(String.format("countryCode = %s", currencyPrediction.getCountryCode())); Currency result = currencyImdg.getSingleObjectBySQL(String.format("countryCode = %s", currencyPrediction.getCountryCode()));
currencyPrediction.setId(result.getId()); currencyPrediction.setId(result.getId());

View file

@ -54,7 +54,7 @@ public class EquitySecurityServiceTest extends AbstractServiceTest {
addRecordToKafka((MockConsumer) equitySecurityService.getConsumer(), TOPIC, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) equitySecurityService.getConsumer(), TOPIC, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); waitingSendAndCheckRecord(ID, mockProducer);
EquitySecurity equityResult = equitySecurityImdg.getSingleObjectBySQL(String.format("fullName = %s", equityPrediction.getFullName())); EquitySecurity equityResult = equitySecurityImdg.getSingleObjectBySQL(String.format("fullName = %s", equityPrediction.getFullName()));
equityPrediction.setId(equityResult.getId()); equityPrediction.setId(equityResult.getId());
@ -94,7 +94,7 @@ public class EquitySecurityServiceTest extends AbstractServiceTest {
addRecordToKafka((MockConsumer) equitySecurityService.getConsumer(), TOPIC, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) equitySecurityService.getConsumer(), TOPIC, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); waitingSendAndCheckRecord(ID, mockProducer);
EquitySecurity equityResult = equitySecurityImdg.getSingleObjectBySQL(String.format("fullName = %s", equityPrediction.getFullName())); EquitySecurity equityResult = equitySecurityImdg.getSingleObjectBySQL(String.format("fullName = %s", equityPrediction.getFullName()));
equityPrediction.setId(equityResult.getId()); equityPrediction.setId(equityResult.getId());
@ -136,7 +136,7 @@ public class EquitySecurityServiceTest extends AbstractServiceTest {
addRecordToKafka((MockConsumer) equitySecurityService.getConsumer(), TOPIC, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) equitySecurityService.getConsumer(), TOPIC, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); waitingSendAndCheckRecord(ID, mockProducer);
EquitySecurity equityResult = equitySecurityImdg.getSingleObjectBySQL(String.format("fullName = %s", equityPrediction.getFullName())); EquitySecurity equityResult = equitySecurityImdg.getSingleObjectBySQL(String.format("fullName = %s", equityPrediction.getFullName()));
equityPrediction.setId(equityResult.getId()); equityPrediction.setId(equityResult.getId());

View file

@ -1,15 +1,11 @@
package ru.spcex.clearing.securities.service; package ru.spcex.clearing.securities.service;
import org.apache.kafka.clients.consumer.MockConsumer; import org.apache.kafka.clients.consumer.MockConsumer;
import org.apache.kafka.clients.producer.MockProducer; import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith; import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.ArgumentCaptor;
import org.mockito.Captor;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.boot.test.mock.mockito.SpyBean;
import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit.jupiter.SpringExtension; import org.springframework.test.context.junit.jupiter.SpringExtension;
import ru.clearing.classes.statics.data.instrument.issue.FixedIncomeCashFlow; import ru.clearing.classes.statics.data.instrument.issue.FixedIncomeCashFlow;
@ -59,10 +55,9 @@ class FixedIncomeCashFlowServiceTest {
private ImdgProvider hazelcastServiceTest; private ImdgProvider hazelcastServiceTest;
private Imdg<FixedIncomeCashFlow> fixedIncomeCashFlowImdg; private Imdg<FixedIncomeCashFlow> fixedIncomeCashFlowImdg;
@Captor @Autowired
private ArgumentCaptor<ProducerRecord> producerRecord; @Qualifier("mockProducer")
@SpyBean protected Producer<String, Object> mockProducer;
private MockProducer<String, Object> mockProducer;
private final String SECURITY_SYMBOL_STR = "777"; private final String SECURITY_SYMBOL_STR = "777";
private final Long SECURITY_SYMBOL_LONG = 777L; private final Long SECURITY_SYMBOL_LONG = 777L;
@ -83,8 +78,9 @@ class FixedIncomeCashFlowServiceTest {
} }
@Test @Test
void fixedIncomeCashFlowNewTest() throws InterruptedException { void fixedIncomeCashFlowNewTest() {
//ARRANGE //ARRANGE
clearAllInImdg(fixedIncomeCashFlowImdg);
FixedIncomeCashFlowNewRequest newRequest = new FixedIncomeCashFlowNewRequest(); FixedIncomeCashFlowNewRequest newRequest = new FixedIncomeCashFlowNewRequest();
newRequest.setSecuritySymbol(SECURITY_SYMBOL_STR); newRequest.setSecuritySymbol(SECURITY_SYMBOL_STR);
newRequest.setNominalValue(TEST_BIG_DECIMAL); newRequest.setNominalValue(TEST_BIG_DECIMAL);
@ -105,14 +101,14 @@ class FixedIncomeCashFlowServiceTest {
addRecordToKafka((MockConsumer) fixedIncomeCashFlowService.getConsumer(), TOPIC_FIXED_INCOME_CASH_FLOW_NEW, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) fixedIncomeCashFlowService.getConsumer(), TOPIC_FIXED_INCOME_CASH_FLOW_NEW, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); waitingSendAndCheckRecord(ID, mockProducer);
FixedIncomeCashFlow resultNew = fixedIncomeCashFlowImdg.getSingleObjectBySQL(String.format("securityId = %s", SECURITY_SYMBOL_STR)); FixedIncomeCashFlow resultNew = fixedIncomeCashFlowImdg.getSingleObjectBySQL(String.format("securityId = %s", SECURITY_SYMBOL_STR));
predictableFixedIncomeCashFlow.setId(resultNew.getId()); predictableFixedIncomeCashFlow.setId(resultNew.getId());
FIXED_INCOME_CASH_FLOW_MATCHER.assertMatch(resultNew, predictableFixedIncomeCashFlow); FIXED_INCOME_CASH_FLOW_MATCHER.assertMatch(resultNew, predictableFixedIncomeCashFlow);
} }
@Test @Test
void fixedIncomeCashFlowUpdateTest() throws InterruptedException { void fixedIncomeCashFlowUpdateTest() {
//ARRANGE //ARRANGE
FixedIncomeCashFlow existFixedIncomeCashFlow = new FixedIncomeCashFlow(); FixedIncomeCashFlow existFixedIncomeCashFlow = new FixedIncomeCashFlow();
existFixedIncomeCashFlow.setSecurityId(SECURITY_SYMBOL_LONG); existFixedIncomeCashFlow.setSecurityId(SECURITY_SYMBOL_LONG);
@ -143,7 +139,7 @@ class FixedIncomeCashFlowServiceTest {
addRecordToKafka((MockConsumer) fixedIncomeCashFlowService.getConsumer(), TOPIC_FIXED_INCOME_CASH_FLOW_UPDATE, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) fixedIncomeCashFlowService.getConsumer(), TOPIC_FIXED_INCOME_CASH_FLOW_UPDATE, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(id, mockProducer, producerRecord); waitingSendAndCheckRecord(id, mockProducer);
FixedIncomeCashFlow resultUpdate = fixedIncomeCashFlowImdg.getSingleObjectBySQL(String.format("securityId = %s", SECURITY_SYMBOL_STR)); FixedIncomeCashFlow resultUpdate = fixedIncomeCashFlowImdg.getSingleObjectBySQL(String.format("securityId = %s", SECURITY_SYMBOL_STR));
predictableFixedIncomeCashFlow.setId(id); predictableFixedIncomeCashFlow.setId(id);
FIXED_INCOME_CASH_FLOW_MATCHER.assertMatch(resultUpdate, predictableFixedIncomeCashFlow); FIXED_INCOME_CASH_FLOW_MATCHER.assertMatch(resultUpdate, predictableFixedIncomeCashFlow);

View file

@ -55,7 +55,7 @@ public class FixedIncomeSecurityServiceTest extends AbstractServiceTest {
addRecordToKafka((MockConsumer) fixedIncomeSecurityService.getConsumer(), TOPIC, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) fixedIncomeSecurityService.getConsumer(), TOPIC, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); waitingSendAndCheckRecord(ID, mockProducer);
FixedIncomeSecurity equityResult = fixedIncomeSecurityImdg.getSingleObjectBySQL(String.format("fullName = %s", fixedIncomePrediction.getFullName())); FixedIncomeSecurity equityResult = fixedIncomeSecurityImdg.getSingleObjectBySQL(String.format("fullName = %s", fixedIncomePrediction.getFullName()));
fixedIncomePrediction.setId(equityResult.getId()); fixedIncomePrediction.setId(equityResult.getId());
@ -96,7 +96,7 @@ public class FixedIncomeSecurityServiceTest extends AbstractServiceTest {
addRecordToKafka((MockConsumer) fixedIncomeSecurityService.getConsumer(), TOPIC, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) fixedIncomeSecurityService.getConsumer(), TOPIC, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); waitingSendAndCheckRecord(ID, mockProducer);
FixedIncomeSecurity equityResult = fixedIncomeSecurityImdg.getSingleObjectBySQL(String.format("fullName = %s", fixedIncomePrediction.getFullName())); FixedIncomeSecurity equityResult = fixedIncomeSecurityImdg.getSingleObjectBySQL(String.format("fullName = %s", fixedIncomePrediction.getFullName()));
fixedIncomePrediction.setId(equityResult.getId()); fixedIncomePrediction.setId(equityResult.getId());
@ -138,7 +138,7 @@ public class FixedIncomeSecurityServiceTest extends AbstractServiceTest {
addRecordToKafka((MockConsumer) fixedIncomeSecurityService.getConsumer(), TOPIC, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) fixedIncomeSecurityService.getConsumer(), TOPIC, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); waitingSendAndCheckRecord(ID, mockProducer);
FixedIncomeSecurity equityResult = fixedIncomeSecurityImdg.getSingleObjectBySQL(String.format("fullName = %s", fixedIncomePrediction.getFullName())); FixedIncomeSecurity equityResult = fixedIncomeSecurityImdg.getSingleObjectBySQL(String.format("fullName = %s", fixedIncomePrediction.getFullName()));
fixedIncomePrediction.setId(equityResult.getId()); fixedIncomePrediction.setId(equityResult.getId());

View file

@ -75,7 +75,7 @@ public class MoneyMarketSecurityServiceTest extends AbstractServiceTest {
addRecordToKafka((MockConsumer) moneyMarketSecurityService.getConsumer(), TOPIC, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) moneyMarketSecurityService.getConsumer(), TOPIC, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); waitingSendAndCheckRecord(ID, mockProducer);
MoneyMarketSecurity moneyMarketSecurityResult = moneyMarketSecurityMap.getSingleObjectBySQL(String.format("fullName = %s", moneyMarketSecurityFactory.getFullName())); MoneyMarketSecurity moneyMarketSecurityResult = moneyMarketSecurityMap.getSingleObjectBySQL(String.format("fullName = %s", moneyMarketSecurityFactory.getFullName()));
moneyMarketSecurityPrediction.setId(moneyMarketSecurityResult.getId()); moneyMarketSecurityPrediction.setId(moneyMarketSecurityResult.getId());
@ -118,7 +118,7 @@ public class MoneyMarketSecurityServiceTest extends AbstractServiceTest {
addRecordToKafka((MockConsumer) moneyMarketSecurityService.getConsumer(), TOPIC, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) moneyMarketSecurityService.getConsumer(), TOPIC, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); waitingSendAndCheckRecord(ID, mockProducer);
MoneyMarketSecurity moneyMarketSecurityResult = moneyMarketSecurityMap.getSingleObjectBySQL(String.format("fullName = %s", moneyMarketSecurityFactory.getFullName())); MoneyMarketSecurity moneyMarketSecurityResult = moneyMarketSecurityMap.getSingleObjectBySQL(String.format("fullName = %s", moneyMarketSecurityFactory.getFullName()));
moneyMarketSecurityPrediction.setId(moneyMarketSecurityResult.getId()); moneyMarketSecurityPrediction.setId(moneyMarketSecurityResult.getId());
@ -167,7 +167,7 @@ public class MoneyMarketSecurityServiceTest extends AbstractServiceTest {
addRecordToKafka((MockConsumer) moneyMarketSecurityService.getConsumer(), TOPIC, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) moneyMarketSecurityService.getConsumer(), TOPIC, PARTITION, 0, jsonString);
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord); waitingSendAndCheckRecord(ID, mockProducer);
MoneyMarketSecurity moneyMarketSecurityResult = moneyMarketSecurityMap.getSingleObjectBySQL(String.format("fullName = %s", moneyMarketSecurityFactory.getFullName())); MoneyMarketSecurity moneyMarketSecurityResult = moneyMarketSecurityMap.getSingleObjectBySQL(String.format("fullName = %s", moneyMarketSecurityFactory.getFullName()));
moneyMarketSecurityPrediction.setId(moneyMarketSecurityResult.getId()); moneyMarketSecurityPrediction.setId(moneyMarketSecurityResult.getId());

View file

@ -68,5 +68,9 @@
<groupId>ru.spcex.platform</groupId> <groupId>ru.spcex.platform</groupId>
<artifactId>platform-enum</artifactId> <artifactId>platform-enum</artifactId>
</dependency> </dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-test</artifactId>
</dependency>
</dependencies> </dependencies>
</project> </project>

View file

@ -4,11 +4,14 @@ import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.ObjectMapper;
import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.MockConsumer; import org.apache.kafka.clients.consumer.MockConsumer;
import org.apache.kafka.clients.producer.MockProducer; import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.clients.producer.RecordMetadata; import org.apache.kafka.clients.producer.RecordMetadata;
import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.TopicPartition;
import org.mockito.ArgumentCaptor; import org.mockito.ArgumentCaptor;
import org.mockito.Mockito;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.requestreply.ReplyingKafkaTemplate;
import ru.spcex.clearing.platform.messaging.domain.ActionType; import ru.spcex.clearing.platform.messaging.domain.ActionType;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest; import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.platform.messaging.domain.Consts; import ru.spcex.clearing.platform.messaging.domain.Consts;
@ -17,6 +20,7 @@ import ru.spcex.clearing.platform.messaging.service.Status;
import ru.spcex.platform.classes.base.SpcexObjectBase; import ru.spcex.platform.classes.base.SpcexObjectBase;
import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.Imdg;
import java.lang.reflect.Field;
import java.util.Collection; import java.util.Collection;
import java.util.Collections; import java.util.Collections;
import java.util.HashMap; import java.util.HashMap;
@ -27,17 +31,50 @@ import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException; import java.util.concurrent.TimeoutException;
import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.mockito.Mockito.timeout; import static org.mockito.Mockito.*;
import static org.mockito.Mockito.verify;
import static ru.spcex.clearing.platform.messaging.service.Status.Success; import static ru.spcex.clearing.platform.messaging.service.Status.Success;
import static ru.spcex.clearing.test.MatcherFactory.usingIgnoringFieldsComparator; import static ru.spcex.clearing.test.MatcherFactory.usingIgnoringFieldsComparator;
import static ru.spcex.clearing.test.config.ImdgTestConfig.defaultAdminId; import static ru.spcex.clearing.test.config.ImdgTestConfig.defaultAdminId;
import static ru.spcex.clearing.test.config.KafkaTestConfig.getCaptor;
public class TestUtils { public class TestUtils {
public static final MatcherFactory.Matcher<BaseRequest<Object>> BASE_REQUEST_MATCHER = usingIgnoringFieldsComparator(); public static final MatcherFactory.Matcher<BaseRequest<Object>> BASE_REQUEST_MATCHER = usingIgnoringFieldsComparator();
private static final ObjectMapper objectMapper = new ObjectMapper(); private static final ObjectMapper objectMapper = new ObjectMapper();
public static void waitingWhenAddedRecordAndCheckIt(Long id, MockProducer mockProducer, ArgumentCaptor<ProducerRecord> producerRecord) { public static void waitingSendAndCheckRecord(Long id, Producer mockProducer, ArgumentCaptor<ProducerRecord> producerRecord) {
//waiting for kafka send message (finale event)
verify(mockProducer, timeout(30_000L).times(1))
.send(producerRecord.capture());
checkRecord(id, producerRecord);
}
public static void waitingSendAndCheckRecord(Long id, Producer mockProducer) {
ArgumentCaptor<ProducerRecord> producerRecord = getCaptor(mockProducer);
//waiting for kafka send message (finale event)
verify(mockProducer, timeout(30_000L).times(1))
.send(producerRecord.capture());
checkRecord(id, producerRecord);
}
public static void waitingSendAndCheckRecord(Long id, KafkaTemplate<String, Object> kafkaTemplate) {
ArgumentCaptor<ProducerRecord> producerRecord = getCaptor(kafkaTemplate);
//waiting for kafka send message (finale event)
verify(kafkaTemplate, timeout(30_000L).times(1))
.send(producerRecord.capture());
checkRecord(id, producerRecord);
}
public static void waitingSendAndReceiveAndCheckRecord(Long id, KafkaTemplate<String, Object> kafkaTemplate) {
ReplyingKafkaTemplate<String, Object, Object> rplKafkaTemplate = (ReplyingKafkaTemplate<String, Object, Object>) kafkaTemplate;
ArgumentCaptor<ProducerRecord> producerRecord = getCaptor(kafkaTemplate);
//waiting for kafka send message (finale event)
verify(rplKafkaTemplate, timeout(30_000L).times(1))
.sendAndReceive(producerRecord.capture());
checkRecord(id, producerRecord);
}
public static void checkRecord(Long id, ArgumentCaptor<ProducerRecord> producerRecord) {
BaseRequest<Object> predictableBaseRequest = new BaseRequest<>(); BaseRequest<Object> predictableBaseRequest = new BaseRequest<>();
predictableBaseRequest.setId(id); predictableBaseRequest.setId(id);
predictableBaseRequest.setActionType(ActionType.SYSTEM); predictableBaseRequest.setActionType(ActionType.SYSTEM);
@ -46,20 +83,16 @@ public class TestUtils {
requestInfoUpdate.setStatus(Success); requestInfoUpdate.setStatus(Success);
predictableBaseRequest.setRequestPayload(requestInfoUpdate); predictableBaseRequest.setRequestPayload(requestInfoUpdate);
//waiting for kafka producer send message (finale event)
verify(mockProducer, timeout(30_000L).times(1))
.send(producerRecord.capture());
BaseRequest<Object> baseRequestResult = (BaseRequest<Object>) producerRecord.getValue().value(); BaseRequest<Object> baseRequestResult = (BaseRequest<Object>) producerRecord.getValue().value();
assertEquals(Consts.REQUEST_INFO_UPDATE, producerRecord.getValue().topic()); assertEquals(Consts.REQUEST_INFO_UPDATE, producerRecord.getValue().topic());
BASE_REQUEST_MATCHER.assertMatch(baseRequestResult, predictableBaseRequest); BASE_REQUEST_MATCHER.assertMatch(baseRequestResult, predictableBaseRequest);
} }
public static void waitingWhenTryAddRecordAndCheckError(Long id, public static void waitingWhenTryAddRecordAndCheckError(Long id,
MockProducer mockProducer, Producer mockProducer,
ArgumentCaptor<ProducerRecord> producerRecord,
String errorCode, String errorCode,
List<String> errorMessageArgs) { List<String> errorMessageArgs) {
ArgumentCaptor<ProducerRecord> producerRecord = getCaptor(mockProducer);
BaseRequest<Object> predictableBaseRequest = new BaseRequest<>(); BaseRequest<Object> predictableBaseRequest = new BaseRequest<>();
predictableBaseRequest.setId(id); predictableBaseRequest.setId(id);
predictableBaseRequest.setActionType(ActionType.SYSTEM); predictableBaseRequest.setActionType(ActionType.SYSTEM);
@ -122,6 +155,21 @@ public class TestUtils {
values.forEach(imdg::delete); values.forEach(imdg::delete);
} }
public static <T, V> T getMockForFildObj(T from, V obj, String fildName) {
T mock = Mockito.mock((Class<T>) from.getClass(), withSettings()
.serializable()
.spiedInstance(from)
.defaultAnswer(CALLS_REAL_METHODS));
try {
Field dbServiceField = obj.getClass().getDeclaredField(fildName);
dbServiceField.setAccessible(true);
dbServiceField.set(obj, mock);
} catch (NoSuchFieldException | IllegalAccessException e) {
throw new RuntimeException(e);
}
return mock;
}
public static class FutureRecordMetadata implements Future<RecordMetadata> { public static class FutureRecordMetadata implements Future<RecordMetadata> {
@Override @Override
public boolean cancel(boolean mayInterruptIfRunning) { public boolean cancel(boolean mayInterruptIfRunning) {

View file

@ -1,31 +1,135 @@
package ru.spcex.clearing.test.config; package ru.spcex.clearing.test.config;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.MockConsumer; import org.apache.kafka.clients.consumer.MockConsumer;
import org.apache.kafka.clients.consumer.OffsetResetStrategy; import org.apache.kafka.clients.consumer.OffsetResetStrategy;
import org.apache.kafka.clients.producer.MockProducer;
import org.apache.kafka.clients.producer.Producer; import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.mockito.ArgumentCaptor;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.beans.factory.config.ConfigurableBeanFactory; import org.springframework.beans.factory.config.ConfigurableBeanFactory;
import org.springframework.boot.test.mock.mockito.MockReset;
import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Import;
import org.springframework.context.annotation.Scope; import org.springframework.context.annotation.Scope;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.requestreply.ReplyingKafkaTemplate;
import org.springframework.kafka.requestreply.RequestReplyFuture;
import org.springframework.kafka.support.SendResult;
import org.springframework.util.concurrent.ListenableFuture;
import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.service.RequestInfo; import ru.spcex.clearing.platform.messaging.service.RequestInfo;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.clearing.test.TestUtils;
import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgId; import ru.spcex.platform.imdg.api.ImdgId;
import ru.spcex.platform.imdg.api.ImdgProvider; import ru.spcex.platform.imdg.api.ImdgProvider;
@Configuration import java.util.HashMap;
public class KafkaTestConfig { import java.util.Map;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import java.util.function.Supplier;
import static org.mockito.Mockito.*;
@Configuration
@Import(ImdgTestConfig.class)
public class KafkaTestConfig {
public final static Map<Producer<String, Object>, ArgumentCaptor<ProducerRecord>> producerCaptors = new HashMap<>();
public final static Map<KafkaTemplate<String, Object>, ArgumentCaptor<ProducerRecord>> templateCaptors = new HashMap<>();
//из-за очисткой перед каждым тестом(MockReset.withSettings(MockReset.AFTER) необходимо каждый раз обновлять doReturn
public static ArgumentCaptor<ProducerRecord> getCaptor(Producer<String, Object> mockProducer){
ArgumentCaptor<ProducerRecord> captor = producerCaptors.get(mockProducer);
setFuture(captor, mockProducer);
return captor;
}
public static ArgumentCaptor<ProducerRecord> getCaptor(KafkaTemplate<String, Object> kafkaTemplate){
ArgumentCaptor<ProducerRecord> captor = templateCaptors.get(kafkaTemplate);
setFuture(captor, (ReplyingKafkaTemplate<String, Object, Object>) kafkaTemplate);
return captor;
}
public static void setFuture(ArgumentCaptor<ProducerRecord> captor, Producer<String, Object> mockProducer){
TestUtils.FutureRecordMetadata future = spy(new TestUtils.FutureRecordMetadata());
doReturn(future).when(mockProducer).send(captor.capture());
}
public static void setFuture(ArgumentCaptor<ProducerRecord> recordArgumentCaptor, ReplyingKafkaTemplate<String, Object, Object> kafkaTemplate){
//sendToQueueWaitForAnswer
RequestReplyFuture<String, Object, Object> replyFuture = spy(RequestReplyFuture.class);
doReturn(replyFuture).when(kafkaTemplate).sendAndReceive(recordArgumentCaptor.capture());
ConsumerRecord<String, Object> consumerRecord = mock(ConsumerRecord.class);
try {
doReturn(consumerRecord).when(replyFuture).get(10, TimeUnit.SECONDS);
} catch (InterruptedException | ExecutionException | TimeoutException | ClassCastException e) {
if (e instanceof InterruptedException) {
Thread.currentThread().interrupt();
}
}
doReturn(null).when(consumerRecord).value();
//sendRequestToQueue
ListenableFuture<SendResult<String, Object>> send = mock(ListenableFuture.class);
doReturn(send).when(kafkaTemplate).send(recordArgumentCaptor.capture());
try {
doReturn(null).when(send).get();
} catch (InterruptedException | ExecutionException e) {
if (e instanceof InterruptedException) {
Thread.currentThread().interrupt();
}
}
}
// @Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE)
@Bean("mockProducer")
public Producer<String, Object> kafkaProducer() {
Producer<String, Object> mockProducer = mock(MockProducer.class, MockReset.withSettings(MockReset.AFTER));
ArgumentCaptor<ProducerRecord> recordArgumentCaptor = ArgumentCaptor.forClass(ProducerRecord.class);
setFuture(recordArgumentCaptor, mockProducer);
producerCaptors.put(mockProducer, recordArgumentCaptor);
return mockProducer;
}
@Autowired
@Bean @Bean
public KafkaSender kafkaSender(Producer<String, Object> kafkaProducer, public KafkaSender kafkaSender(ImdgProvider imdgProvider, Producer<String, Object> mockProducer) {
ImdgProvider imdgProvider) {
ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator(); ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator();
return KafkaSender return KafkaSender
.setup() .setup()
.producer(kafkaProducer) .producer(mockProducer)
.idGenerator(imdgIdGenerator::nextId)
.imdgProvider(s -> {
Imdg<RequestInfo> imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class);
return imdg::insert;
})
.build();
}
@Bean("kafkaTestTemplate")
public KafkaTemplate<String, Object> kafkaTemplate() {
ReplyingKafkaTemplate<String, Object, Object> kafkaTemplate = mock(ReplyingKafkaTemplate.class, MockReset.withSettings(MockReset.AFTER));
ArgumentCaptor<ProducerRecord> recordArgumentCaptor = ArgumentCaptor.forClass(ProducerRecord.class);
setFuture(recordArgumentCaptor, kafkaTemplate);
templateCaptors.put(kafkaTemplate, recordArgumentCaptor);
return kafkaTemplate;
}
@Autowired
@Bean
public Supplier<KafkaSender> kafkaSenderSupplier(KafkaTemplate<String, Object> kafkaTemplate,
@Qualifier("hazelcastServiceTest") ImdgProvider imdgProvider) {
ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator();
return () -> KafkaSender
.setup()
.setKafkaTemplate(kafkaTemplate)
.idGenerator(imdgIdGenerator::nextId) .idGenerator(imdgIdGenerator::nextId)
.imdgProvider(s -> { .imdgProvider(s -> {
Imdg<RequestInfo> imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class); Imdg<RequestInfo> imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class);

View file

@ -40,7 +40,7 @@ public class TradeImporterService {
private final Supplier<KafkaSender> kafka; private final Supplier<KafkaSender> kafka;
private final IMessageResolver messageResolver; private final IMessageResolver messageResolver;
@Value("${trade-importer.database.schema}") @Value("${trade-importer.database.schema:SPVB_TS}")
private String schema; private String schema;

View file

@ -1,20 +1,13 @@
package ru.spcex.clearing.trade.importer; package ru.spcex.clearing.trade.importer;
import org.apache.kafka.clients.producer.MockProducer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.junit.jupiter.api.extension.ExtendWith; 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.Autowired;
import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.boot.test.mock.mockito.MockBean; import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit.jupiter.SpringExtension; import org.springframework.test.context.junit.jupiter.SpringExtension;
import ru.clearing.classes.statics.data.misc.STrades; import ru.clearing.classes.statics.data.misc.STrades;
import ru.clearing.classes.statics.data.scheduler.PlannerAllToday;
import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.imdg.IMDGDistributedNames;
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.ImdgTestConfig;
import ru.spcex.clearing.test.config.KafkaTestConfig; import ru.spcex.clearing.test.config.KafkaTestConfig;
import ru.spcex.clearing.trade.importer.config.ErrorResolverConfig; import ru.spcex.clearing.trade.importer.config.ErrorResolverConfig;
@ -25,9 +18,6 @@ import ru.spcex.clearing.trade.importer.services.TradeImporterService;
import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider; import ru.spcex.platform.imdg.api.ImdgProvider;
import static org.mockito.Mockito.doReturn;
import static org.mockito.Mockito.spy;
import static ru.spcex.clearing.test.MatcherFactory.usingIgnoringFieldsComparator;
import static ru.spcex.clearing.test.config.ImdgTestConfig.currentID; import static ru.spcex.clearing.test.config.ImdgTestConfig.currentID;
import static ru.spcex.clearing.test.config.ImdgTestConfig.waitAvailableImdgProviderAndAddAdminWithDefaultId; import static ru.spcex.clearing.test.config.ImdgTestConfig.waitAvailableImdgProviderAndAddAdminWithDefaultId;
@ -41,14 +31,13 @@ import static ru.spcex.clearing.test.config.ImdgTestConfig.waitAvailableImdgProv
ImdgTestConfig.class, ImdgTestConfig.class,
KafkaTestConfig.class}) KafkaTestConfig.class})
public abstract class AbstractServiceTest { public abstract class AbstractServiceTest {
protected static final MatcherFactory.Matcher<PlannerAllToday> PLANNER_ALL_TODAY_MATCHER = usingIgnoringFieldsComparator("created", "updated");
protected static final long id = currentID.getAndIncrement(); protected static final long id = currentID.getAndIncrement();
protected Imdg<STrades> sTradesImdg; protected Imdg<STrades> sTradesImdg;
@Captor @Autowired
protected ArgumentCaptor<ProducerRecord> producerRecord; @Qualifier("kafkaTestTemplate")
@MockBean protected KafkaTemplate<String, Object> kafkaTemplate;
protected MockProducer<String, Object> mockProducer;
@Autowired @Autowired
@Qualifier("hazelcastServiceTest") @Qualifier("hazelcastServiceTest")
protected ImdgProvider imdgProvider; protected ImdgProvider imdgProvider;
@ -56,8 +45,5 @@ public abstract class AbstractServiceTest {
protected void init() { protected void init() {
waitAvailableImdgProviderAndAddAdminWithDefaultId(); waitAvailableImdgProviderAndAddAdminWithDefaultId();
this.sTradesImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_STrades, STrades.class); this.sTradesImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_STrades, STrades.class);
TestUtils.FutureRecordMetadata future = spy(new TestUtils.FutureRecordMetadata());
doReturn(future).when(mockProducer).send(producerRecord.capture());
} }
} }

View file

@ -20,7 +20,7 @@ public class DbTestConnectionConfig {
public SingleConnectionDataSource dataSource() { public SingleConnectionDataSource dataSource() {
String login = "sa"; String login = "sa";
String password = "Aa123456"; String password = "Aa123456";
String dbUrl = "jdbc:sqlserver://localhost:1433;database=SPVB_TS;schema=dbo"; String dbUrl = "jdbc:sqlserver://10.200.200.144:1433;database=ni";
SingleConnectionDataSource cpds = new SingleConnectionDataSource(); SingleConnectionDataSource cpds = new SingleConnectionDataSource();

View file

@ -1,9 +1,11 @@
package ru.spcex.clearing.trade.importer.services; package ru.spcex.clearing.trade.importer.services;
import org.apache.kafka.clients.consumer.MockConsumer; import org.apache.kafka.clients.consumer.MockConsumer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Test;
import org.mockito.ArgumentCaptor;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import ru.clearing.classes.statics.data.misc.STrades; import ru.clearing.classes.statics.data.misc.STrades;
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest; import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
@ -22,6 +24,7 @@ import static org.mockito.Mockito.timeout;
import static org.mockito.Mockito.verify; import static org.mockito.Mockito.verify;
import static ru.spcex.clearing.test.MatcherFactory.usingIgnoringFieldsComparator; import static ru.spcex.clearing.test.MatcherFactory.usingIgnoringFieldsComparator;
import static ru.spcex.clearing.test.TestUtils.*; import static ru.spcex.clearing.test.TestUtils.*;
import static ru.spcex.clearing.test.config.KafkaTestConfig.getCaptor;
class TradeImporterServiceTest extends AbstractServiceTest { class TradeImporterServiceTest extends AbstractServiceTest {
public static final MatcherFactory.Matcher<STrades> S_TRADES_MATCHER = usingIgnoringFieldsComparator(); public static final MatcherFactory.Matcher<STrades> S_TRADES_MATCHER = usingIgnoringFieldsComparator();
@ -63,8 +66,9 @@ class TradeImporterServiceTest extends AbstractServiceTest {
addRecordToKafka((MockConsumer) launcherCommandReceiver.getConsumer(), Task.getOfTrades.topic(), 0, 1, getJsonStringForNew(new LauncherCommandRequest(),0)); addRecordToKafka((MockConsumer) launcherCommandReceiver.getConsumer(), Task.getOfTrades.topic(), 0, 1, getJsonStringForNew(new LauncherCommandRequest(),0));
//waiting for kafka producer send message //waiting for kafka producer send message
verify(mockProducer, timeout(30_000L).times(1)) ArgumentCaptor<ProducerRecord> captor = getCaptor(kafkaTemplate);
.send(producerRecord.capture()); verify(kafkaTemplate, timeout(30_000L).times(1))
.send(captor.capture());
STrades sTrades = tradeImporterService.getSTradesFromImdg(sTrade, sTradesImdg); STrades sTrades = tradeImporterService.getSTradesFromImdg(sTrade, sTradesImdg);