Simple refactoring company-service tests.

This commit is contained in:
psemenkov 2023-02-22 18:56:31 +03:00
parent 3acebda42b
commit a1076c5026
10 changed files with 364 additions and 371 deletions

View file

@ -44,7 +44,6 @@ import static org.mockito.Mockito.verify;
import static ru.spcex.clearing.account.utils.MatcherFactory.usingIgnoringFieldsComparator; import static ru.spcex.clearing.account.utils.MatcherFactory.usingIgnoringFieldsComparator;
import static ru.spcex.clearing.account.utils.TestUtils.*; import static ru.spcex.clearing.account.utils.TestUtils.*;
import static ru.spcex.clearing.platform.messaging.service.Status.Error; import static ru.spcex.clearing.platform.messaging.service.Status.Error;
import static ru.spcex.clearing.platform.messaging.service.Status.Success;
@ExtendWith(SpringExtension.class) @ExtendWith(SpringExtension.class)
@ContextConfiguration(classes = { @ContextConfiguration(classes = {
@ -116,14 +115,6 @@ public class BankAccountServiceTest {
@Test @Test
public void bankAccountNew() { public void bankAccountNew() {
//ARRANGE //ARRANGE
BaseRequest<Object> predictableBaseRequest = new BaseRequest<>();
predictableBaseRequest.setId(ID);
predictableBaseRequest.setActionType(ActionType.SYSTEM);
RequestInfoUpdate requestInfoUpdate = new RequestInfoUpdate();
requestInfoUpdate.setId(ID);
requestInfoUpdate.setStatus(Success);
predictableBaseRequest.setRequestPayload(requestInfoUpdate);
BankAccount predictableBankAccount = getBankAccount(); BankAccount predictableBankAccount = getBankAccount();
Company company = getTestCompany(); Company company = getTestCompany();
companyImdg.insert(company); companyImdg.insert(company);
@ -134,11 +125,7 @@ public class BankAccountServiceTest {
//ACT //ACT
addRecordToKafka((MockConsumer) bankAccountService.getConsumer(), TOPIC_ACCOUNT_NEW, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) bankAccountService.getConsumer(), TOPIC_ACCOUNT_NEW, PARTITION, 0, jsonString);
waitingWhenAddedRecordAndCheckIt(ID, producer, producerRecord);
//waiting for kafka producer send message (finale event)
verify(producer, timeout(30_000L).times(1))
.send(producerRecord.capture());
BaseRequest<Object> baseRequestResult = (BaseRequest<Object>) producerRecord.getValue().value();
Account accountResult = accountImdg.getSingleObjectBySQL(String.format("account = %s", acc)); Account accountResult = accountImdg.getSingleObjectBySQL(String.format("account = %s", acc));
BankAccount bankAccountResult = bankAccountImdg.getSingleObjectBySQL(String.format("account = %s or companyId = %s", acc, addresseeIdNew)); BankAccount bankAccountResult = bankAccountImdg.getSingleObjectBySQL(String.format("account = %s or companyId = %s", acc, addresseeIdNew));
@ -147,10 +134,8 @@ public class BankAccountServiceTest {
setSameValueToField(accountResult, predictableAccount); setSameValueToField(accountResult, predictableAccount);
//ASSERT //ASSERT
assertEquals(Consts.REQUEST_INFO_UPDATE, producerRecord.getValue().topic());
BANK_ACCOUNT_MATCHER.assertMatch(bankAccountResult, predictableBankAccount); BANK_ACCOUNT_MATCHER.assertMatch(bankAccountResult, predictableBankAccount);
ACCOUNT_MATCHER.assertMatch(accountResult, predictableAccount); ACCOUNT_MATCHER.assertMatch(accountResult, predictableAccount);
BASE_REQUEST_MATCHER.assertMatch(baseRequestResult, predictableBaseRequest);
} }
/** /**
@ -285,14 +270,6 @@ public class BankAccountServiceTest {
@Test @Test
void bankAccountUpdate() { void bankAccountUpdate() {
//ARRANGE //ARRANGE
BaseRequest<Object> predictableBaseRequest = new BaseRequest<>();
predictableBaseRequest.setId(ID);
predictableBaseRequest.setActionType(ActionType.SYSTEM);
RequestInfoUpdate requestInfoUpdate = new RequestInfoUpdate();
requestInfoUpdate.setId(ID);
requestInfoUpdate.setStatus(Success);
predictableBaseRequest.setRequestPayload(requestInfoUpdate);
// Company company = getTestCompany(); // Company company = getTestCompany();
// companyImdg.insert(company); // companyImdg.insert(company);
Account predictableAccount = getTestAccount(accountId, acc); Account predictableAccount = getTestAccount(accountId, acc);
@ -329,10 +306,8 @@ public class BankAccountServiceTest {
//ACT //ACT
addRecordToKafka((MockConsumer) bankAccountService.getConsumer(), TOPIC_ACCOUNT_UPDATE, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) bankAccountService.getConsumer(), TOPIC_ACCOUNT_UPDATE, PARTITION, 0, jsonString);
//waiting for kafka producer send message (finale event) //ASSERT
verify(producer, timeout(30_000L).times(1)) waitingWhenAddedRecordAndCheckIt(ID, producer, producerRecord);
.send(producerRecord.capture());
BaseRequest<Object> baseRequestResult = (BaseRequest<Object>) producerRecord.getValue().value();
Account accountResult = accountImdg.getSingleObjectBySQL(String.format("account = %s", acc)); Account accountResult = accountImdg.getSingleObjectBySQL(String.format("account = %s", acc));
BankAccount resultUpdating = bankAccountImdg.getSingleObjectBySQL(String.format("account = %s or companyId = %s", acc, addresseeIdNew)); BankAccount resultUpdating = bankAccountImdg.getSingleObjectBySQL(String.format("account = %s or companyId = %s", acc, addresseeIdNew));
@ -340,14 +315,8 @@ public class BankAccountServiceTest {
predictableUpdateBankAccount.setCompanyId(resultUpdating.getCompanyId()); predictableUpdateBankAccount.setCompanyId(resultUpdating.getCompanyId());
predictableAccount.setUpdated(accountResult.getUpdated()); predictableAccount.setUpdated(accountResult.getUpdated());
//ASSERT
assertEquals(Consts.REQUEST_INFO_UPDATE, producerRecord.getValue().topic());
//ASSERT
assertEquals(Consts.REQUEST_INFO_UPDATE, producerRecord.getValue().topic());
BANK_ACCOUNT_MATCHER.assertMatch(resultUpdating, predictableUpdateBankAccount); BANK_ACCOUNT_MATCHER.assertMatch(resultUpdating, predictableUpdateBankAccount);
ACCOUNT_MATCHER.assertMatch(accountResult, predictableAccount); ACCOUNT_MATCHER.assertMatch(accountResult, predictableAccount);
BASE_REQUEST_MATCHER.assertMatch(baseRequestResult, predictableBaseRequest);
} }
/** /**
@ -372,14 +341,11 @@ public class BankAccountServiceTest {
//ACT //ACT
addRecordToKafka((MockConsumer) bankAccountService.getConsumer(), TOPIC_ACCOUNT_DELETE, PARTITION, 0, jsonString); addRecordToKafka((MockConsumer) bankAccountService.getConsumer(), TOPIC_ACCOUNT_DELETE, PARTITION, 0, jsonString);
//waiting for kafka producer send message (finale event)
verify(producer, timeout(30_000L).times(1))
.send(producerRecord.capture());
//ASSERT //ASSERT
waitingWhenAddedRecordAndCheckIt(ID, producer, producerRecord);
BankAccount bankAccount = bankAccountImdg.getSingleObjectByID(ID); BankAccount bankAccount = bankAccountImdg.getSingleObjectByID(ID);
Assertions.assertNull(bankAccount); Assertions.assertNull(bankAccount);
assertEquals(Consts.REQUEST_INFO_UPDATE, producerRecord.getValue().topic());
} }
private Company getTestCompany() { private Company getTestCompany() {

View file

@ -4,10 +4,15 @@ 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.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 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.service.RequestInfoUpdate;
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;
@ -19,9 +24,34 @@ import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit; 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.mockito.Mockito.timeout;
import static org.mockito.Mockito.verify;
import static ru.spcex.clearing.account.utils.MatcherFactory.usingIgnoringFieldsComparator;
import static ru.spcex.clearing.platform.messaging.service.Status.Success;
public class TestUtils { public class TestUtils {
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) {
BaseRequest<Object> predictableBaseRequest = new BaseRequest<>();
predictableBaseRequest.setId(id);
predictableBaseRequest.setActionType(ActionType.SYSTEM);
RequestInfoUpdate requestInfoUpdate = new RequestInfoUpdate();
requestInfoUpdate.setId(id);
requestInfoUpdate.setStatus(Success);
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();
assertEquals(Consts.REQUEST_INFO_UPDATE, producerRecord.getValue().topic());
BASE_REQUEST_MATCHER.assertMatch(baseRequestResult, predictableBaseRequest);
}
public static void addRecordToKafka(MockConsumer mockConsumer, String topic, int partition, long offset, String jsonValue) { public static void addRecordToKafka(MockConsumer mockConsumer, String topic, int partition, long offset, String jsonValue) {
TopicPartition tp = new TopicPartition(topic, partition); TopicPartition tp = new TopicPartition(topic, partition);
HashMap<TopicPartition, Long> startOffsets = new HashMap<>(); HashMap<TopicPartition, Long> startOffsets = new HashMap<>();

View file

@ -43,7 +43,7 @@ public class HazelcastServiceTestConfiguration {
joinConfig.setTcpIpConfig(new TcpIpConfig().setEnabled(true).setMembers(List.of("127.0.0.1"))); joinConfig.setTcpIpConfig(new TcpIpConfig().setEnabled(true).setMembers(List.of("127.0.0.1")));
networkConfig.setJoin(joinConfig); networkConfig.setJoin(joinConfig);
cfg.setNetworkConfig(networkConfig); cfg.setNetworkConfig(networkConfig);
hazelcastInstance = Hazelcast.newHazelcastInstance(cfg); hazelcastInstance = Hazelcast.getOrCreateHazelcastInstance(cfg);
HazelcastHelper.otcSystem_setStorageState(true, hazelcastInstance); HazelcastHelper.otcSystem_setStorageState(true, hazelcastInstance);
return new HazelcastService(taskExecutorHazelcastClientInitializer, taskExecutorIdGeneratorAwaiter, params); return new HazelcastService(taskExecutorHazelcastClientInitializer, taskExecutorIdGeneratorAwaiter, params);
} }

View file

@ -0,0 +1,44 @@
package ru.spcex.clearing.company.config;
import org.apache.kafka.clients.consumer.MockConsumer;
import org.apache.kafka.clients.consumer.OffsetResetStrategy;
import org.apache.kafka.clients.producer.Producer;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.beans.factory.config.ConfigurableBeanFactory;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Scope;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.service.RequestInfo;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgId;
import ru.spcex.platform.imdg.api.ImdgProvider;
@Configuration
public class KafkaConfigTest {
@Autowired
@Bean(name = "kafkaSenderTest")
public KafkaSender kafkaSender(Producer<String, Object> kafkaProducer, @Qualifier("hazelcastServiceTest") ImdgProvider imdgProvider) {
ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator();
return KafkaSender
.setup()
.producer(kafkaProducer)
.idGenerator(imdgIdGenerator::nextId)
.imdgProvider(s -> {
Imdg<RequestInfo> imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class);
return imdg::insert;
})
.build();
}
@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE)
@Bean(name = "mockConsumerTest")
public MockConsumer<String, Object> createConsumer() {
return new MockConsumer<>(OffsetResetStrategy.EARLIEST);
}
}

View file

@ -1,42 +1,41 @@
package ru.spcex.clearing.company.service; package ru.spcex.clearing.company.service;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.hazelcast.core.IMap;
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.producer.MockProducer; import org.apache.kafka.clients.producer.MockProducer;
import org.apache.kafka.common.TopicPartition; import org.apache.kafka.clients.producer.ProducerRecord;
import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.BeforeEach;
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.generated.ClearingMemberCategory; import ru.clearing.classes.statics.data.generated.ClearingMemberCategory;
import ru.spcex.clearing.company.config.HazelcastServiceTestConfiguration; import ru.spcex.clearing.company.config.HazelcastServiceTestConfiguration;
import ru.spcex.clearing.company.utils.ImapEvent; import ru.spcex.clearing.company.config.KafkaConfigTest;
import ru.spcex.clearing.company.utils.MatcherFactory.Matcher; import ru.spcex.clearing.company.utils.MatcherFactory.Matcher;
import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.ActionType;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest; import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.platform.messaging.domain.Consts; import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteRequest; import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.company.ClearingMemberCategoryNewRequest; import ru.spcex.clearing.platform.messaging.domain.cud.company.ClearingMemberCategoryNewRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.company.ClearingMemberCategoryUpdateRequest; import ru.spcex.clearing.platform.messaging.domain.cud.company.ClearingMemberCategoryUpdateRequest;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService; import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService;
import java.util.Collections; import javax.annotation.PostConstruct;
import java.util.HashMap;
import static ru.spcex.clearing.company.config.HazelcastServiceTestConfiguration.currentID; import static ru.spcex.clearing.company.config.HazelcastServiceTestConfiguration.currentID;
import static ru.spcex.clearing.company.utils.MatcherFactory.usingIgnoringFieldsComparator; import static ru.spcex.clearing.company.utils.MatcherFactory.usingIgnoringFieldsComparator;
import static ru.spcex.clearing.company.utils.TestUtils.*;
@ExtendWith(SpringExtension.class) @ExtendWith(SpringExtension.class)
@ContextConfiguration(classes = { @ContextConfiguration(classes = {
ClearingMemberCategoryService.class,
KafkaConfigTest.class,
HazelcastServiceTestConfiguration.class}) HazelcastServiceTestConfiguration.class})
class ClearingMemberCategoryServiceTest { class ClearingMemberCategoryServiceTest {
@ -46,67 +45,44 @@ class ClearingMemberCategoryServiceTest {
private static final String TOPIC_MEMBER_CATEGORY_UPDATE = Consts.DESTINATION_CLEARING_MEMBER_CATEGORY_UPDATE; private static final String TOPIC_MEMBER_CATEGORY_UPDATE = Consts.DESTINATION_CLEARING_MEMBER_CATEGORY_UPDATE;
private static final String TOPIC_MEMBER_CATEGORY_DELETE = Consts.DESTINATION_CLEARING_MEMBER_CATEGORY_DELETE; private static final String TOPIC_MEMBER_CATEGORY_DELETE = Consts.DESTINATION_CLEARING_MEMBER_CATEGORY_DELETE;
private static final Long ID = currentID.getAndIncrement(); private static final Long ID = currentID.getAndIncrement();
@Autowired
ClearingMemberCategoryService clearingMemberCategoryService;
@Autowired @Autowired
@Qualifier("hazelcastServiceTest") @Qualifier("hazelcastServiceTest")
private HazelcastService hazelcastServiceTest; private HazelcastService hazelcastServiceTest;
private MockConsumer<String, Object> mockConsumer; private Imdg<ClearingMemberCategory> memberCategoryImdg;
@Captor
private ArgumentCaptor<ProducerRecord> producerRecord;
@SpyBean
private MockProducer<String, Object> mockProducer; private MockProducer<String, Object> mockProducer;
@BeforeEach @PostConstruct
void setUp() { private void init() {
mockConsumer = new MockConsumer<>(OffsetResetStrategy.EARLIEST); hazelcastServiceTest.waitAvailable();
mockProducer = new MockProducer<>(); memberCategoryImdg = hazelcastServiceTest.getImdg(IMDGDistributedNames.Map_ClearingMemberCategory, ClearingMemberCategory.class);
} }
@Test @Test
void clearingMemberCategoryNew() throws InterruptedException { void clearingMemberCategoryNew() throws InterruptedException {
//ARRANGE //ARRANGE
String clearingMemberCategory = "1234";
ClearingMemberCategoryNewRequest memberCategoryNewRequest = new ClearingMemberCategoryNewRequest(); ClearingMemberCategoryNewRequest memberCategoryNewRequest = new ClearingMemberCategoryNewRequest();
memberCategoryNewRequest.setClearingMemberCategory("1234"); memberCategoryNewRequest.setClearingMemberCategory(clearingMemberCategory);
BaseRequest<ClearingMemberCategoryNewRequest> baseUpdateRequest = new BaseRequest<>();
baseUpdateRequest.setRequestPayload(memberCategoryNewRequest);
baseUpdateRequest.setId(ID);
baseUpdateRequest.setActionType(ActionType.NEW);
String jsonBaseForNewRequest;
ObjectMapper objectMapper = new ObjectMapper();
try {
jsonBaseForNewRequest = objectMapper.writeValueAsString(baseUpdateRequest);
} catch (JsonProcessingException e) {
throw new RuntimeException(e);
}
ClearingMemberCategory predictableClearingMemberCategory = new ClearingMemberCategory(); ClearingMemberCategory predictableClearingMemberCategory = new ClearingMemberCategory();
predictableClearingMemberCategory.setId(ID); predictableClearingMemberCategory.setId(ID);
predictableClearingMemberCategory.setClearingMemberCategory("1234"); predictableClearingMemberCategory.setClearingMemberCategory("1234");
//ACT //ACT
//service set up String jsonString = getJsonStringForNew(memberCategoryNewRequest, ID);
ClearingMemberCategoryService clearingMemberCategoryService = new ClearingMemberCategoryService(mockConsumer, mockProducer, hazelcastServiceTest);
//callbacks set up addRecordToKafka((MockConsumer) clearingMemberCategoryService.getConsumer(), TOPIC_MEMBER_CATEGORY_NEW, PARTITION, 0, jsonString);
clearingMemberCategoryService.afterPropertiesSet();
IMap<Long, ClearingMemberCategory> iMap = hazelcastServiceTest.getHazelcast().getMap(IMDGDistributedNames.Map_ClearingMemberCategory);
ImapEvent imapEvent = new ImapEvent(iMap);
//KAFKA
mockConsumer.schedulePollTask(() -> {
mockConsumer.rebalance(Collections.singletonList(new TopicPartition(TOPIC_MEMBER_CATEGORY_NEW, PARTITION)));
mockConsumer.addRecord(new ConsumerRecord<>(TOPIC_MEMBER_CATEGORY_NEW, PARTITION, 0, "key", jsonBaseForNewRequest));
});
HashMap<TopicPartition, Long> startOffsetsUpdating = new HashMap<>();
TopicPartition tpUpdating = new TopicPartition(TOPIC_MEMBER_CATEGORY_NEW, PARTITION);
startOffsetsUpdating.put(tpUpdating, 0L);
mockConsumer.updateBeginningOffsets(startOffsetsUpdating);
//waiting for hazelcast map item updates
imapEvent.waitWhenHappened();
//ASSERT //ASSERT
ClearingMemberCategory resultUpdating = iMap.get(ID); waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord);
MEMBER_CATEGORY_MATCHER.assertMatch(resultUpdating, predictableClearingMemberCategory); ClearingMemberCategory resultNew = memberCategoryImdg.getSingleObjectBySQL(String.format("clearingMemberCategory = %s", clearingMemberCategory));
MEMBER_CATEGORY_MATCHER.assertMatch(resultNew, predictableClearingMemberCategory);
} }
/** /**
@ -119,55 +95,28 @@ class ClearingMemberCategoryServiceTest {
@Test @Test
void clearingMemberCategoryUpdate() throws InterruptedException { void clearingMemberCategoryUpdate() throws InterruptedException {
//ARRANGE //ARRANGE
String clearingMemberCategory = "1234";
ClearingMemberCategory existsСlearingMemberCategory = new ClearingMemberCategory(); ClearingMemberCategory existsСlearingMemberCategory = new ClearingMemberCategory();
existsСlearingMemberCategory.setId(ID); existsСlearingMemberCategory.setId(ID);
existsСlearingMemberCategory.setClearingMemberCategory("0000"); existsСlearingMemberCategory.setClearingMemberCategory("0000");
memberCategoryImdg.insert(existsСlearingMemberCategory);
ClearingMemberCategoryUpdateRequest memberCategoryUpdateRequest = new ClearingMemberCategoryUpdateRequest(); ClearingMemberCategoryUpdateRequest memberCategoryUpdateRequest = new ClearingMemberCategoryUpdateRequest();
memberCategoryUpdateRequest.setId(ID); memberCategoryUpdateRequest.setId(ID);
memberCategoryUpdateRequest.setClearingMemberCategory("1234"); memberCategoryUpdateRequest.setClearingMemberCategory(clearingMemberCategory);
BaseRequest<ClearingMemberCategoryUpdateRequest> baseUpdateRequest = new BaseRequest<>();
baseUpdateRequest.setRequestPayload(memberCategoryUpdateRequest);
baseUpdateRequest.setId(ID);
baseUpdateRequest.setActionType(ActionType.UPDATE);
String jsonBaseForUpdatingRequest;
ObjectMapper objectMapper = new ObjectMapper();
try {
jsonBaseForUpdatingRequest = objectMapper.writeValueAsString(baseUpdateRequest);
} catch (JsonProcessingException e) {
throw new RuntimeException(e);
}
ClearingMemberCategory predictableClearingMemberCategory = new ClearingMemberCategory(); ClearingMemberCategory predictableClearingMemberCategory = new ClearingMemberCategory();
predictableClearingMemberCategory.setId(ID); predictableClearingMemberCategory.setId(ID);
predictableClearingMemberCategory.setClearingMemberCategory("1234"); predictableClearingMemberCategory.setClearingMemberCategory(clearingMemberCategory);
//ACT //ACT
//service set up String jsonString = getJsonStringForUPDATE(memberCategoryUpdateRequest, ID);
ClearingMemberCategoryService clearingMemberCategoryService = new ClearingMemberCategoryService(mockConsumer, mockProducer, hazelcastServiceTest); addRecordToKafka((MockConsumer) clearingMemberCategoryService.getConsumer(), TOPIC_MEMBER_CATEGORY_UPDATE, PARTITION, 0, jsonString);
//callbacks set up
clearingMemberCategoryService.afterPropertiesSet();
IMap<Long, ClearingMemberCategory> iMap = hazelcastServiceTest.getHazelcast().getMap(IMDGDistributedNames.Map_ClearingMemberCategory);
iMap.put(ID, existsСlearingMemberCategory);
ImapEvent imapEvent = new ImapEvent(iMap);
//KAFKA
mockConsumer.schedulePollTask(() -> {
mockConsumer.rebalance(Collections.singletonList(new TopicPartition(TOPIC_MEMBER_CATEGORY_UPDATE, PARTITION)));
mockConsumer.addRecord(new ConsumerRecord<>(TOPIC_MEMBER_CATEGORY_UPDATE, PARTITION, 0, "key", jsonBaseForUpdatingRequest));
});
HashMap<TopicPartition, Long> startOffsetsUpdating = new HashMap<>();
TopicPartition tpUpdating = new TopicPartition(TOPIC_MEMBER_CATEGORY_UPDATE, PARTITION);
startOffsetsUpdating.put(tpUpdating, 0L);
mockConsumer.updateBeginningOffsets(startOffsetsUpdating);
//waiting for hazelcast map item updates
imapEvent.waitWhenHappened();
//ASSERT //ASSERT
ClearingMemberCategory resultUpdating = iMap.get(ID); waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord);
ClearingMemberCategory resultUpdating = memberCategoryImdg.getSingleObjectBySQL(String.format("clearingMemberCategory = %s", clearingMemberCategory));
MEMBER_CATEGORY_MATCHER.assertMatch(resultUpdating, predictableClearingMemberCategory); MEMBER_CATEGORY_MATCHER.assertMatch(resultUpdating, predictableClearingMemberCategory);
} }
@ -180,50 +129,23 @@ class ClearingMemberCategoryServiceTest {
@Test @Test
void clearingMemberCategoryDelete() throws InterruptedException { void clearingMemberCategoryDelete() throws InterruptedException {
//ARRANGE //ARRANGE
String clearingMemberCategory = "0000";
ClearingMemberCategory existsСlearingMemberCategory = new ClearingMemberCategory(); ClearingMemberCategory existsСlearingMemberCategory = new ClearingMemberCategory();
existsСlearingMemberCategory.setId(ID); existsСlearingMemberCategory.setId(ID);
existsСlearingMemberCategory.setClearingMemberCategory("0000"); existsСlearingMemberCategory.setClearingMemberCategory(clearingMemberCategory);
memberCategoryImdg.insert(existsСlearingMemberCategory);
CommonDeleteRequest memberCategoryDeleteRequest = new CommonDeleteRequest(); CommonDeleteRequest memberCategoryDeleteRequest = new CommonDeleteRequest();
memberCategoryDeleteRequest.setId(ID); memberCategoryDeleteRequest.setId(ID);
BaseRequest<CommonDeleteRequest> baseDeleteRequest = new BaseRequest<>();
baseDeleteRequest.setRequestPayload(memberCategoryDeleteRequest);
baseDeleteRequest.setId(ID);
baseDeleteRequest.setActionType(ActionType.DELETE);
String jsonBaseForDeleteRequest;
ObjectMapper objectMapper = new ObjectMapper();
try {
jsonBaseForDeleteRequest = objectMapper.writeValueAsString(baseDeleteRequest);
} catch (JsonProcessingException e) {
throw new RuntimeException(e);
}
//ACT //ACT
//service set up String jsonString = getJsonStringForDELETE(memberCategoryDeleteRequest, ID);
ClearingMemberCategoryService clearingMemberCategoryService = new ClearingMemberCategoryService(mockConsumer, mockProducer, hazelcastServiceTest); addRecordToKafka((MockConsumer) clearingMemberCategoryService.getConsumer(), TOPIC_MEMBER_CATEGORY_DELETE, PARTITION, 0, jsonString);
//callbacks set up
clearingMemberCategoryService.afterPropertiesSet();
IMap<Long, ClearingMemberCategory> iMap = hazelcastServiceTest.getHazelcast().getMap(IMDGDistributedNames.Map_ClearingMemberCategory);
iMap.put(ID, existsСlearingMemberCategory);
ImapEvent imapEvent = new ImapEvent(iMap);
//KAFKA
mockConsumer.schedulePollTask(() -> {
mockConsumer.rebalance(Collections.singletonList(new TopicPartition(TOPIC_MEMBER_CATEGORY_DELETE, PARTITION)));
mockConsumer.addRecord(new ConsumerRecord<>(TOPIC_MEMBER_CATEGORY_DELETE, PARTITION, 0, "key", jsonBaseForDeleteRequest));
});
HashMap<TopicPartition, Long> startOffsetsUpdating = new HashMap<>();
TopicPartition tpDeleting = new TopicPartition(TOPIC_MEMBER_CATEGORY_DELETE, PARTITION);
startOffsetsUpdating.put(tpDeleting, 0L);
mockConsumer.updateBeginningOffsets(startOffsetsUpdating);
//waiting for hazelcast map item removes
imapEvent.waitWhenHappened();
//ASSERT //ASSERT
Assertions.assertEquals(0, iMap.size()); waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord);
ClearingMemberCategory resultDeleting = memberCategoryImdg.getSingleObjectBySQL(String.format("clearingMemberCategory = %s", clearingMemberCategory));
Assertions.assertNull(resultDeleting);
} }
} }

View file

@ -1,39 +1,38 @@
package ru.spcex.clearing.company.service; package ru.spcex.clearing.company.service;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.hazelcast.core.IMap;
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.producer.MockProducer; import org.apache.kafka.clients.producer.MockProducer;
import org.apache.kafka.common.TopicPartition; import org.apache.kafka.clients.producer.ProducerRecord;
import org.junit.jupiter.api.BeforeEach;
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;
import ru.clearing.classes.statics.data.profile.CompanyInfo; import ru.clearing.classes.statics.data.profile.CompanyInfo;
import ru.spcex.clearing.company.config.HazelcastServiceTestConfiguration; import ru.spcex.clearing.company.config.HazelcastServiceTestConfiguration;
import ru.spcex.clearing.company.utils.ImapEvent; import ru.spcex.clearing.company.config.KafkaConfigTest;
import ru.spcex.clearing.company.utils.MatcherFactory; import ru.spcex.clearing.company.utils.MatcherFactory;
import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.ActionType;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest; import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.platform.messaging.domain.Consts; import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.company.CompanyInfoUpdateRequest; import ru.spcex.clearing.platform.messaging.domain.cud.company.CompanyInfoUpdateRequest;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService; import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService;
import java.util.Collections; import javax.annotation.PostConstruct;
import java.util.HashMap;
import static ru.spcex.clearing.company.utils.MatcherFactory.usingIgnoringFieldsComparator; import static ru.spcex.clearing.company.utils.MatcherFactory.usingIgnoringFieldsComparator;
import static ru.spcex.clearing.company.utils.TestUtils.*;
@ExtendWith(SpringExtension.class) @ExtendWith(SpringExtension.class)
@ContextConfiguration(classes = { @ContextConfiguration(classes = {
CompanyInfoService.class,
KafkaConfigTest.class,
HazelcastServiceTestConfiguration.class}) HazelcastServiceTestConfiguration.class})
class CompanyInfoServiceTest { class CompanyInfoServiceTest {
@ -41,17 +40,22 @@ class CompanyInfoServiceTest {
private static final int PARTITION = 0; private static final int PARTITION = 0;
private static final String TOPIC_COMPANY_INFO_UPDATE = Consts.DESTINATION_COMPANY_INFO_UPDATE; private static final String TOPIC_COMPANY_INFO_UPDATE = Consts.DESTINATION_COMPANY_INFO_UPDATE;
private static final Long ID = 0L; private static final Long ID = 0L;
@Autowired
CompanyInfoService companyInfoService;
@Autowired @Autowired
@Qualifier("hazelcastServiceTest") @Qualifier("hazelcastServiceTest")
private HazelcastService hazelcastServiceTest; private HazelcastService hazelcastServiceTest;
private MockConsumer<String, Object> mockConsumer; private Imdg<Company> companyImdg;
@Captor
private ArgumentCaptor<ProducerRecord> producerRecord;
@SpyBean
private MockProducer<String, Object> mockProducer; private MockProducer<String, Object> mockProducer;
@BeforeEach @PostConstruct
void setUp() { private void init() {
mockConsumer = new MockConsumer<>(OffsetResetStrategy.EARLIEST); hazelcastServiceTest.waitAvailable();
mockProducer = new MockProducer<>(); companyImdg = hazelcastServiceTest.getImdg(IMDGDistributedNames.Map_Company, Company.class);
} }
/** /**
@ -91,6 +95,7 @@ class CompanyInfoServiceTest {
Company existsCompany = new Company(); Company existsCompany = new Company();
existsCompany.setId(ID); existsCompany.setId(ID);
existsCompany.setProfile(existsCompanyInfo); existsCompany.setProfile(existsCompanyInfo);
companyImdg.insert(existsCompany);
CompanyInfoUpdateRequest companyInfoUpdateRequest = new CompanyInfoUpdateRequest(); CompanyInfoUpdateRequest companyInfoUpdateRequest = new CompanyInfoUpdateRequest();
companyInfoUpdateRequest.setId(ID); companyInfoUpdateRequest.setId(ID);
@ -106,18 +111,6 @@ class CompanyInfoServiceTest {
companyInfoUpdateRequest.setShortName("updated shortName"); companyInfoUpdateRequest.setShortName("updated shortName");
companyInfoUpdateRequest.setFullName("updated fullName"); companyInfoUpdateRequest.setFullName("updated fullName");
BaseRequest<CompanyInfoUpdateRequest> baseUpdateRequest = new BaseRequest<>();
baseUpdateRequest.setRequestPayload(companyInfoUpdateRequest);
baseUpdateRequest.setId(ID);
baseUpdateRequest.setActionType(ActionType.UPDATE);
String jsonBaseForUpdatingRequest;
ObjectMapper objectMapper = new ObjectMapper();
try {
jsonBaseForUpdatingRequest = objectMapper.writeValueAsString(baseUpdateRequest);
} catch (JsonProcessingException e) {
throw new RuntimeException(e);
}
CompanyInfo predictableCompanyInfo = new CompanyInfo(); CompanyInfo predictableCompanyInfo = new CompanyInfo();
predictableCompanyInfo.setId(ID); predictableCompanyInfo.setId(ID);
predictableCompanyInfo.setCompanyId(ID); predictableCompanyInfo.setCompanyId(ID);
@ -134,31 +127,14 @@ class CompanyInfoServiceTest {
predictableCompanyInfo.setFullName("updated fullName"); predictableCompanyInfo.setFullName("updated fullName");
//ACT //ACT
//service set up String jsonString = getJsonStringForUPDATE(companyInfoUpdateRequest, ID);
CompanyInfoService clearingMemberCategoryService = new CompanyInfoService(mockConsumer, mockProducer, hazelcastServiceTest);
//callbacks set up addRecordToKafka((MockConsumer) companyInfoService.getConsumer(), TOPIC_COMPANY_INFO_UPDATE, PARTITION, 0, jsonString);
clearingMemberCategoryService.afterPropertiesSet();
IMap<Long, Company> iMap = hazelcastServiceTest.getHazelcast().getMap(IMDGDistributedNames.Map_Company);
iMap.put(ID, existsCompany);
ImapEvent imapEvent = new ImapEvent(iMap);
//KAFKA
mockConsumer.schedulePollTask(() -> {
mockConsumer.rebalance(Collections.singletonList(new TopicPartition(TOPIC_COMPANY_INFO_UPDATE, PARTITION)));
mockConsumer.addRecord(new ConsumerRecord<>(TOPIC_COMPANY_INFO_UPDATE, PARTITION, 0, "key", jsonBaseForUpdatingRequest));
});
HashMap<TopicPartition, Long> startOffsetsUpdating = new HashMap<>();
TopicPartition tpUpdating = new TopicPartition(TOPIC_COMPANY_INFO_UPDATE, PARTITION);
startOffsetsUpdating.put(tpUpdating, 0L);
mockConsumer.updateBeginningOffsets(startOffsetsUpdating);
//waiting for hazelcast map updates
imapEvent.waitWhenHappened();
//ASSERT //ASSERT
Company resultUpdating = iMap.get(ID); waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord);
Company resultUpdating = companyImdg.getSingleObjectByID(ID);
COMPANY_INFO_MATCHER.assertMatch(resultUpdating.getProfile(), predictableCompanyInfo); COMPANY_INFO_MATCHER.assertMatch(resultUpdating.getProfile(), predictableCompanyInfo);
} }
} }

View file

@ -1,40 +1,39 @@
package ru.spcex.clearing.company.service; package ru.spcex.clearing.company.service;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.hazelcast.core.IMap;
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.producer.MockProducer; import org.apache.kafka.clients.producer.MockProducer;
import org.apache.kafka.common.TopicPartition; import org.apache.kafka.clients.producer.ProducerRecord;
import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.BeforeEach;
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;
import ru.clearing.classes.statics.data.generated.ClearingMemberCategory; import ru.clearing.classes.statics.data.generated.ClearingMemberCategory;
import ru.spcex.clearing.company.config.HazelcastServiceTestConfiguration; import ru.spcex.clearing.company.config.HazelcastServiceTestConfiguration;
import ru.spcex.clearing.company.utils.ImapEvent; import ru.spcex.clearing.company.config.KafkaConfigTest;
import ru.spcex.clearing.company.utils.MatcherFactory; import ru.spcex.clearing.company.utils.MatcherFactory;
import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.ActionType;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest; import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.platform.messaging.domain.Consts; import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteRequest; import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteRequest;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService; import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService;
import java.util.Collections; import javax.annotation.PostConstruct;
import java.util.HashMap;
import static ru.spcex.clearing.company.utils.MatcherFactory.usingIgnoringFieldsComparator; import static ru.spcex.clearing.company.utils.MatcherFactory.usingIgnoringFieldsComparator;
import static ru.spcex.clearing.company.utils.TestUtils.*;
@ExtendWith(SpringExtension.class) @ExtendWith(SpringExtension.class)
@ContextConfiguration(classes = { @ContextConfiguration(classes = {
CompanyService.class,
KafkaConfigTest.class,
HazelcastServiceTestConfiguration.class}) HazelcastServiceTestConfiguration.class})
class CompanyServiceTest { class CompanyServiceTest {
@ -43,16 +42,22 @@ class CompanyServiceTest {
private static final String TOPIC_COMPANY_DELETE = Consts.DESTINATION_COMPANY_DELETE; private static final String TOPIC_COMPANY_DELETE = Consts.DESTINATION_COMPANY_DELETE;
private static final Long ID = 0L; private static final Long ID = 0L;
@Autowired
CompanyService companyService;
@Autowired @Autowired
@Qualifier("hazelcastServiceTest") @Qualifier("hazelcastServiceTest")
private HazelcastService hazelcastServiceTest; private HazelcastService hazelcastServiceTest;
private MockConsumer<String, Object> mockConsumer; private Imdg<Company> companyImdg;
@Captor
private ArgumentCaptor<ProducerRecord> producerRecord;
@SpyBean
private MockProducer<String, Object> mockProducer; private MockProducer<String, Object> mockProducer;
@BeforeEach @PostConstruct
void setUp() { private void init() {
mockConsumer = new MockConsumer<>(OffsetResetStrategy.EARLIEST); hazelcastServiceTest.waitAvailable();
mockProducer = new MockProducer<>(); companyImdg = hazelcastServiceTest.getImdg(IMDGDistributedNames.Map_Company, Company.class);
} }
/** /**
@ -66,47 +71,19 @@ class CompanyServiceTest {
//ARRANGE //ARRANGE
Company existsCompany = new Company(); Company existsCompany = new Company();
existsCompany.setId(ID); existsCompany.setId(ID);
companyImdg.insert(existsCompany);
CommonDeleteRequest memberCategoryDeleteRequest = new CommonDeleteRequest(); CommonDeleteRequest commonDeleteRequest = new CommonDeleteRequest();
memberCategoryDeleteRequest.setId(ID); commonDeleteRequest.setId(ID);
BaseRequest<CommonDeleteRequest> baseDeleteRequest = new BaseRequest<>();
baseDeleteRequest.setRequestPayload(memberCategoryDeleteRequest);
baseDeleteRequest.setId(ID);
baseDeleteRequest.setActionType(ActionType.DELETE);
String jsonBaseForDeleteRequest;
ObjectMapper objectMapper = new ObjectMapper();
try {
jsonBaseForDeleteRequest = objectMapper.writeValueAsString(baseDeleteRequest);
} catch (JsonProcessingException e) {
throw new RuntimeException(e);
}
//ACT //ACT
//service set up String jsonString = getJsonStringForDELETE(commonDeleteRequest, ID);
CompanyService companyService = new CompanyService(mockConsumer, mockProducer, hazelcastServiceTest); addRecordToKafka((MockConsumer) companyService.getConsumer(), TOPIC_COMPANY_DELETE, PARTITION, 0, jsonString);
//callbacks set up
companyService.afterPropertiesSet();
IMap<Long, Company> iMap = hazelcastServiceTest.getHazelcast().getMap(IMDGDistributedNames.Map_Company);
iMap.put(ID, existsCompany);
ImapEvent imapEvent = new ImapEvent(iMap);
//KAFKA
mockConsumer.schedulePollTask(() -> {
mockConsumer.rebalance(Collections.singletonList(new TopicPartition(TOPIC_COMPANY_DELETE, PARTITION)));
mockConsumer.addRecord(new ConsumerRecord<>(TOPIC_COMPANY_DELETE, PARTITION, 0, "key", jsonBaseForDeleteRequest));
});
HashMap<TopicPartition, Long> startOffsetsUpdating = new HashMap<>();
TopicPartition tpDeleting = new TopicPartition(TOPIC_COMPANY_DELETE, PARTITION);
startOffsetsUpdating.put(tpDeleting, 0L);
mockConsumer.updateBeginningOffsets(startOffsetsUpdating);
//waiting for hazelcast map item removes
imapEvent.waitWhenHappened();
//ASSERT //ASSERT
Assertions.assertEquals(0, iMap.size()); waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord);
Company resultDeleting = companyImdg.getSingleObjectByID(ID);
Assertions.assertNull(resultDeleting);
} }
} }

View file

@ -1,38 +1,37 @@
package ru.spcex.clearing.company.service; package ru.spcex.clearing.company.service;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.hazelcast.core.IMap;
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.producer.MockProducer; import org.apache.kafka.clients.producer.MockProducer;
import org.apache.kafka.common.TopicPartition; import org.apache.kafka.clients.producer.ProducerRecord;
import org.junit.jupiter.api.BeforeEach;
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.CompanySymbols; import ru.clearing.classes.statics.data.company.CompanySymbols;
import ru.spcex.clearing.company.config.HazelcastServiceTestConfiguration; import ru.spcex.clearing.company.config.HazelcastServiceTestConfiguration;
import ru.spcex.clearing.company.utils.ImapEvent; import ru.spcex.clearing.company.config.KafkaConfigTest;
import ru.spcex.clearing.company.utils.MatcherFactory; import ru.spcex.clearing.company.utils.MatcherFactory;
import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.ActionType;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest; import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.platform.messaging.domain.Consts; import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.company.CompanySymbolUpdateRequest; import ru.spcex.clearing.platform.messaging.domain.cud.company.CompanySymbolUpdateRequest;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService; import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService;
import java.util.Collections; import javax.annotation.PostConstruct;
import java.util.HashMap;
import static ru.spcex.clearing.company.utils.MatcherFactory.usingIgnoringFieldsComparator; import static ru.spcex.clearing.company.utils.MatcherFactory.usingIgnoringFieldsComparator;
import static ru.spcex.clearing.company.utils.TestUtils.*;
@ExtendWith(SpringExtension.class) @ExtendWith(SpringExtension.class)
@ContextConfiguration(classes = { @ContextConfiguration(classes = {
CompanySymbolService.class,
KafkaConfigTest.class,
HazelcastServiceTestConfiguration.class}) HazelcastServiceTestConfiguration.class})
class CompanySymbolServiceTest { class CompanySymbolServiceTest {
@ -41,16 +40,22 @@ class CompanySymbolServiceTest {
private static final String TOPIC_COMPANY_SYMBOL_UPDATE = Consts.DESTINATION_COMPANY_SYMBOL_UPDATE; private static final String TOPIC_COMPANY_SYMBOL_UPDATE = Consts.DESTINATION_COMPANY_SYMBOL_UPDATE;
private static final Long ID = 0L; private static final Long ID = 0L;
@Autowired
CompanySymbolService companySymbolService;
@Autowired @Autowired
@Qualifier("hazelcastServiceTest") @Qualifier("hazelcastServiceTest")
private HazelcastService hazelcastServiceTest; private HazelcastService hazelcastServiceTest;
private MockConsumer<String, Object> mockConsumer; private Imdg<CompanySymbols> companySymbolsImdg;
@Captor
private ArgumentCaptor<ProducerRecord> producerRecord;
@SpyBean
private MockProducer<String, Object> mockProducer; private MockProducer<String, Object> mockProducer;
@BeforeEach @PostConstruct
void setUp() { private void init() {
mockConsumer = new MockConsumer<>(OffsetResetStrategy.EARLIEST); hazelcastServiceTest.waitAvailable();
mockProducer = new MockProducer<>(); companySymbolsImdg = hazelcastServiceTest.getImdg(IMDGDistributedNames.Map_CompanySymbols, CompanySymbols.class);
} }
/** /**
@ -65,62 +70,34 @@ class CompanySymbolServiceTest {
@Test @Test
void companySymbolUpdate() throws InterruptedException { void companySymbolUpdate() throws InterruptedException {
//ARRANGE //ARRANGE
String companySymbolValue = "newCompanySymbolValue";
CompanySymbols existsCompanySymbols = new CompanySymbols(); CompanySymbols existsCompanySymbols = new CompanySymbols();
existsCompanySymbols.setId(ID); existsCompanySymbols.setId(ID);
existsCompanySymbols.setCompanyId(ID); existsCompanySymbols.setCompanyId(ID);
existsCompanySymbols.setCompanySymbol("0000"); existsCompanySymbols.setCompanySymbol("0000");
existsCompanySymbols.setCompanySymbolValue("exists companySymbolValue"); existsCompanySymbols.setCompanySymbolValue("exists companySymbolValue");
companySymbolsImdg.insert(existsCompanySymbols);
CompanySymbolUpdateRequest companySymbolUpdateRequest = new CompanySymbolUpdateRequest(); CompanySymbolUpdateRequest companySymbolUpdateRequest = new CompanySymbolUpdateRequest();
companySymbolUpdateRequest.setId(ID); companySymbolUpdateRequest.setId(ID);
companySymbolUpdateRequest.setCompanyId(ID); companySymbolUpdateRequest.setCompanyId(ID);
companySymbolUpdateRequest.setCompanySymbol("1234"); companySymbolUpdateRequest.setCompanySymbol("1234");
companySymbolUpdateRequest.setCompanySymbolValue("new companySymbolValue"); companySymbolUpdateRequest.setCompanySymbolValue(companySymbolValue);
BaseRequest<CompanySymbolUpdateRequest> baseUpdateRequest = new BaseRequest<>();
baseUpdateRequest.setRequestPayload(companySymbolUpdateRequest);
baseUpdateRequest.setId(ID);
baseUpdateRequest.setActionType(ActionType.UPDATE);
String jsonBaseForUpdatingRequest;
ObjectMapper objectMapper = new ObjectMapper();
try {
jsonBaseForUpdatingRequest = objectMapper.writeValueAsString(baseUpdateRequest);
} catch (JsonProcessingException e) {
throw new RuntimeException(e);
}
CompanySymbols predictableCompanySymbols = new CompanySymbols(); CompanySymbols predictableCompanySymbols = new CompanySymbols();
predictableCompanySymbols.setId(ID); predictableCompanySymbols.setId(ID);
predictableCompanySymbols.setCompanyId(ID); predictableCompanySymbols.setCompanyId(ID);
predictableCompanySymbols.setCompanySymbol("0000"); predictableCompanySymbols.setCompanySymbol("0000");
predictableCompanySymbols.setCompanySymbolValue("new companySymbolValue"); predictableCompanySymbols.setCompanySymbolValue(companySymbolValue);
//ACT //ACT
//service set up String jsonString = getJsonStringForUPDATE(companySymbolUpdateRequest, ID);
CompanySymbolService companySymbolService = new CompanySymbolService(mockConsumer, mockProducer, hazelcastServiceTest); addRecordToKafka((MockConsumer) companySymbolService.getConsumer(), TOPIC_COMPANY_SYMBOL_UPDATE, PARTITION, 0, jsonString);
//callbacks set up
companySymbolService.afterPropertiesSet();
IMap<Long, CompanySymbols> iMap = hazelcastServiceTest.getHazelcast().getMap(IMDGDistributedNames.Map_CompanySymbols);
iMap.put(ID, existsCompanySymbols);
ImapEvent imapEvent = new ImapEvent(iMap);
//KAFKA
mockConsumer.schedulePollTask(() -> {
mockConsumer.rebalance(Collections.singletonList(new TopicPartition(TOPIC_COMPANY_SYMBOL_UPDATE, PARTITION)));
mockConsumer.addRecord(new ConsumerRecord<>(TOPIC_COMPANY_SYMBOL_UPDATE, PARTITION, 0, "key", jsonBaseForUpdatingRequest));
});
HashMap<TopicPartition, Long> startOffsetsUpdating = new HashMap<>();
TopicPartition tpUpdating = new TopicPartition(TOPIC_COMPANY_SYMBOL_UPDATE, PARTITION);
startOffsetsUpdating.put(tpUpdating, 0L);
mockConsumer.updateBeginningOffsets(startOffsetsUpdating);
//waiting for hazelcast map item updates
imapEvent.waitWhenHappened();
//ASSERT //ASSERT
CompanySymbols resultUpdating = iMap.get(ID); waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord);
CompanySymbols resultUpdating = companySymbolsImdg.getSingleObjectBySQL(String.format("companySymbolValue = %s", companySymbolValue));
COMPANY_SYMBOL_MATCHER.assertMatch(resultUpdating, predictableCompanySymbols); COMPANY_SYMBOL_MATCHER.assertMatch(resultUpdating, predictableCompanySymbols);
} }
} }

View file

@ -1,38 +1,37 @@
package ru.spcex.clearing.company.service; package ru.spcex.clearing.company.service;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.hazelcast.core.IMap;
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.producer.MockProducer; import org.apache.kafka.clients.producer.MockProducer;
import org.apache.kafka.common.TopicPartition; import org.apache.kafka.clients.producer.ProducerRecord;
import org.junit.jupiter.api.BeforeEach;
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.profile.Contact; import ru.clearing.classes.statics.data.profile.Contact;
import ru.spcex.clearing.company.config.HazelcastServiceTestConfiguration; import ru.spcex.clearing.company.config.HazelcastServiceTestConfiguration;
import ru.spcex.clearing.company.utils.ImapEvent; import ru.spcex.clearing.company.config.KafkaConfigTest;
import ru.spcex.clearing.company.utils.MatcherFactory; import ru.spcex.clearing.company.utils.MatcherFactory;
import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.ActionType;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest; import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.platform.messaging.domain.Consts; import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.company.ContactUpdateRequest; import ru.spcex.clearing.platform.messaging.domain.cud.company.ContactUpdateRequest;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService; import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService;
import java.util.Collections; import javax.annotation.PostConstruct;
import java.util.HashMap;
import static ru.spcex.clearing.company.utils.MatcherFactory.usingIgnoringFieldsComparator; import static ru.spcex.clearing.company.utils.MatcherFactory.usingIgnoringFieldsComparator;
import static ru.spcex.clearing.company.utils.TestUtils.*;
@ExtendWith(SpringExtension.class) @ExtendWith(SpringExtension.class)
@ContextConfiguration(classes = { @ContextConfiguration(classes = {
ContactService.class,
KafkaConfigTest.class,
HazelcastServiceTestConfiguration.class}) HazelcastServiceTestConfiguration.class})
class ContactServiceTest { class ContactServiceTest {
@ -41,16 +40,23 @@ class ContactServiceTest {
private static final String TOPIC_CONTACT_UPDATE = Consts.DESTINATION_CONTACT_UPDATE; private static final String TOPIC_CONTACT_UPDATE = Consts.DESTINATION_CONTACT_UPDATE;
private static final Long ID = 0L; private static final Long ID = 0L;
@Autowired
ContactService contactService;
@Autowired @Autowired
@Qualifier("hazelcastServiceTest") @Qualifier("hazelcastServiceTest")
private HazelcastService hazelcastServiceTest; private HazelcastService hazelcastServiceTest;
private MockConsumer<String, Object> mockConsumer; private Imdg<Contact> contactImdg;
@Captor
private ArgumentCaptor<ProducerRecord> producerRecord;
@SpyBean
private MockProducer<String, Object> mockProducer; private MockProducer<String, Object> mockProducer;
@BeforeEach @PostConstruct
void setUp() { private void init() {
mockConsumer = new MockConsumer<>(OffsetResetStrategy.EARLIEST); hazelcastServiceTest.waitAvailable();
mockProducer = new MockProducer<>(); contactImdg = hazelcastServiceTest.getImdg(IMDGDistributedNames.Map_Contact, Contact.class);
} }
/** /**
@ -70,6 +76,7 @@ class ContactServiceTest {
existsContact.setCompanyId(ID); existsContact.setCompanyId(ID);
existsContact.setContactType("0000"); existsContact.setContactType("0000");
existsContact.setContactValue("exists ContactValue"); existsContact.setContactValue("exists ContactValue");
contactImdg.insert(existsContact);
ContactUpdateRequest contactUpdateRequest = new ContactUpdateRequest(); ContactUpdateRequest contactUpdateRequest = new ContactUpdateRequest();
contactUpdateRequest.setId(ID); contactUpdateRequest.setId(ID);
@ -77,18 +84,6 @@ class ContactServiceTest {
contactUpdateRequest.setContactType("1234"); contactUpdateRequest.setContactType("1234");
contactUpdateRequest.setContactValue("new ContactValue"); contactUpdateRequest.setContactValue("new ContactValue");
BaseRequest<ContactUpdateRequest> baseUpdateRequest = new BaseRequest<>();
baseUpdateRequest.setRequestPayload(contactUpdateRequest);
baseUpdateRequest.setId(ID);
baseUpdateRequest.setActionType(ActionType.UPDATE);
String jsonBaseForUpdatingRequest;
ObjectMapper objectMapper = new ObjectMapper();
try {
jsonBaseForUpdatingRequest = objectMapper.writeValueAsString(baseUpdateRequest);
} catch (JsonProcessingException e) {
throw new RuntimeException(e);
}
Contact predictableContact = new Contact(); Contact predictableContact = new Contact();
predictableContact.setId(ID); predictableContact.setId(ID);
predictableContact.setCompanyId(ID); predictableContact.setCompanyId(ID);
@ -96,30 +91,13 @@ class ContactServiceTest {
predictableContact.setContactValue("new ContactValue"); predictableContact.setContactValue("new ContactValue");
//ACT //ACT
//service set up String jsonString = getJsonStringForUPDATE(contactUpdateRequest, ID);
ContactService contactService = new ContactService(mockConsumer, mockProducer, hazelcastServiceTest); addRecordToKafka((MockConsumer) contactService.getConsumer(), TOPIC_CONTACT_UPDATE, PARTITION, 0, jsonString);
contactService.afterPropertiesSet();
//callbacks set up
IMap<Long, Contact> iMap = hazelcastServiceTest.getHazelcast().getMap(IMDGDistributedNames.Map_Contact);
iMap.put(ID, existsContact);
ImapEvent imapEvent = new ImapEvent(iMap);
//KAFKA
mockConsumer.schedulePollTask(() -> {
mockConsumer.rebalance(Collections.singletonList(new TopicPartition(TOPIC_CONTACT_UPDATE, PARTITION)));
mockConsumer.addRecord(new ConsumerRecord<>(TOPIC_CONTACT_UPDATE, PARTITION, 0, "key", jsonBaseForUpdatingRequest));
});
HashMap<TopicPartition, Long> startOffsetsUpdating = new HashMap<>();
TopicPartition tpUpdating = new TopicPartition(TOPIC_CONTACT_UPDATE, PARTITION);
startOffsetsUpdating.put(tpUpdating, 0L);
mockConsumer.updateBeginningOffsets(startOffsetsUpdating);
//waiting for hazelcast map updates
imapEvent.waitWhenHappened();
//ASSERT //ASSERT
Contact resultUpdating = iMap.get(ID); waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord);
Contact resultUpdating = contactImdg.getSingleObjectByID(ID);
CONTACT_MATCHER.assertMatch(resultUpdating, predictableContact); CONTACT_MATCHER.assertMatch(resultUpdating, predictableContact);
} }
} }

View file

@ -0,0 +1,123 @@
package ru.spcex.clearing.company.utils;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.MockConsumer;
import org.apache.kafka.clients.producer.MockProducer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.clients.producer.RecordMetadata;
import org.apache.kafka.common.TopicPartition;
import org.mockito.ArgumentCaptor;
import ru.spcex.clearing.platform.messaging.domain.ActionType;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate;
import ru.spcex.platform.classes.base.SpcexObjectBase;
import ru.spcex.platform.imdg.api.Imdg;
import java.util.Collection;
import java.util.Collections;
import java.util.HashMap;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.mockito.Mockito.timeout;
import static org.mockito.Mockito.verify;
import static ru.spcex.clearing.company.utils.MatcherFactory.usingIgnoringFieldsComparator;
import static ru.spcex.clearing.platform.messaging.service.Status.Success;
public class TestUtils {
public static final MatcherFactory.Matcher<BaseRequest<Object>> BASE_REQUEST_MATCHER = usingIgnoringFieldsComparator();
private static final ObjectMapper objectMapper = new ObjectMapper();
public static void waitingWhenAddedRecordAndCheckIt(Long id, MockProducer mockProducer, ArgumentCaptor<ProducerRecord> producerRecord) {
BaseRequest<Object> predictableBaseRequest = new BaseRequest<>();
predictableBaseRequest.setId(id);
predictableBaseRequest.setActionType(ActionType.SYSTEM);
RequestInfoUpdate requestInfoUpdate = new RequestInfoUpdate();
requestInfoUpdate.setId(id);
requestInfoUpdate.setStatus(Success);
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();
assertEquals(Consts.REQUEST_INFO_UPDATE, producerRecord.getValue().topic());
BASE_REQUEST_MATCHER.assertMatch(baseRequestResult, predictableBaseRequest);
}
public static void addRecordToKafka(MockConsumer mockConsumer, String topic, int partition, long offset, String jsonValue) {
TopicPartition tp = new TopicPartition(topic, partition);
HashMap<TopicPartition, Long> startOffsets = new HashMap<>();
startOffsets.put(tp, 0L);
mockConsumer.updateBeginningOffsets(startOffsets);
mockConsumer.schedulePollTask(() -> {
mockConsumer.rebalance(Collections.singletonList(tp));
mockConsumer.addRecord(new ConsumerRecord<>(topic, partition, offset, "key", jsonValue));
});
}
public static <T> String getJsonStringForNew(T accountRequest, long id) {
return getJsonBaseRequest(accountRequest, id, ActionType.NEW);
}
public static <T> String getJsonStringForUPDATE(T accountRequest, long id) {
return getJsonBaseRequest(accountRequest, id, ActionType.UPDATE);
}
public static <T> String getJsonStringForDELETE(T accountRequest, long id) {
return getJsonBaseRequest(accountRequest, id, ActionType.DELETE);
}
private static <T> String getJsonBaseRequest(T accountRequest, long id, ActionType actionType) {
BaseRequest<T> baseRequest = new BaseRequest<>();
baseRequest.setRequestPayload(accountRequest);
baseRequest.setId(id);
baseRequest.setActionType(actionType);
String jsonBaseRequest;
try {
jsonBaseRequest = objectMapper.writeValueAsString(baseRequest);
} catch (JsonProcessingException e) {
throw new RuntimeException(e);
}
return jsonBaseRequest;
}
public static <T extends SpcexObjectBase> void clearAllInImdg(Imdg<T> imdg) {
Collection<T> values = imdg.getAllValues();
values.forEach(imdg::delete);
}
public static class FutureRecordMetadata implements Future<RecordMetadata> {
@Override
public boolean cancel(boolean mayInterruptIfRunning) {
return false;
}
@Override
public boolean isCancelled() {
return false;
}
@Override
public boolean isDone() {
return false;
}
@Override
public RecordMetadata get() throws InterruptedException, ExecutionException {
return null;
}
@Override
public RecordMetadata get(long timeout, TimeUnit unit) throws InterruptedException, ExecutionException, TimeoutException {
return null;
}
}
}