Correction of the error in the service-account and the service-company has been completed. Both service are successfully launched in maven.
This commit is contained in:
parent
614b785344
commit
4fc3b9b4f5
10 changed files with 164 additions and 141 deletions
|
|
@ -59,6 +59,16 @@ class AccountServiceTest {
|
|||
@Qualifier("mockConsumerTest")
|
||||
private MockConsumer<String, Object> mockConsumer;
|
||||
|
||||
/**
|
||||
* {@link AccountService#accountNew(BaseRequest)}<br>
|
||||
* Тест проверяет создание сущности {@link BaseRequest} в Hazelcast при передаче из Apache Kafka.<br>
|
||||
* Входной запрос {@link AccountSdf01Request}:<br>
|
||||
* {@link AccountSdfRequestPart#setSdfId} - текущий Id<br>
|
||||
* {@link AccountSdfRequestPart#setAccount} - 123456789123<br>
|
||||
* {@link AccountSdfRequestPart#setCompanyId} - текущий Id<br>
|
||||
* {@link AccountSdf01Request#setGroupingSdf01Id} - текущий Id<br>
|
||||
* {@link AccountSdf01Request#setAccounts} - Collections.singletonList(AccountSdfRequestPart)<br>
|
||||
*/
|
||||
@Test
|
||||
void accountNew() throws InterruptedException {
|
||||
//arrange
|
||||
|
|
@ -129,7 +139,7 @@ class AccountServiceTest {
|
|||
IMap<Long, RequestInfo> requestInfoIMap = hazelcastServiceTest.getHazelcast().getMap(IMDGDistributedNames.Map_RequestInfo);
|
||||
|
||||
ImapEvent imapEvent = new ImapEvent(requestInfoIMap);
|
||||
imapEvent.waitHappened();
|
||||
imapEvent.waitWhenHappened();
|
||||
//ASSERT
|
||||
Account accountResult = accountIMap.get(firstID);
|
||||
RequestInfo requestInfoResult = requestInfoIMap.get(secondID);
|
||||
|
|
|
|||
|
|
@ -139,7 +139,7 @@ public class BankAccountServiceTest {
|
|||
|
||||
//waiting for hazelcast map item updates
|
||||
ImapEvent imapEvent = new ImapEvent(iMap);
|
||||
imapEvent.waitHappened();
|
||||
imapEvent.waitWhenHappened();
|
||||
|
||||
//ASSERT
|
||||
BankAccount result = iMap.get(ID);
|
||||
|
|
@ -235,7 +235,7 @@ public class BankAccountServiceTest {
|
|||
mockConsumer.updateBeginningOffsets(startOffsetsUpdating);
|
||||
|
||||
//waiting for hazelcast map item updates
|
||||
imapEvent.waitHappened();
|
||||
imapEvent.waitWhenHappened();
|
||||
|
||||
//ASSERT
|
||||
BankAccount resultUpdating = iMap.get(ID);
|
||||
|
|
@ -301,7 +301,7 @@ public class BankAccountServiceTest {
|
|||
mockConsumer.updateBeginningOffsets(startOffsetsDeleting);
|
||||
|
||||
//waiting for hazelcast map item updates
|
||||
imapEvent.waitHappened();
|
||||
imapEvent.waitWhenHappened();
|
||||
|
||||
//ASSERT
|
||||
Assertions.assertEquals(0, iMap.size());
|
||||
|
|
|
|||
|
|
@ -38,13 +38,16 @@ public class ImapEvent<T> {
|
|||
}, false);
|
||||
}
|
||||
|
||||
public void waitHappened() throws InterruptedException {
|
||||
public void waitWhenHappened() throws InterruptedException {
|
||||
//running timer task as daemon thread
|
||||
Timer timer = new Timer(true);
|
||||
timer.scheduleAtFixedRate(new TimerTask() {
|
||||
boolean secondRan;
|
||||
|
||||
@Override
|
||||
public void run() {
|
||||
checkEventHappened.set(true);//если что-то пойдет не так не тормозить основной поток
|
||||
checkEventHappened.set(secondRan);//если что-то пойдет не так не тормозить основной поток
|
||||
secondRan = true;
|
||||
}
|
||||
}, 0, 60 * 1000);
|
||||
synchronized (checkEventHappened) {
|
||||
|
|
|
|||
|
|
@ -13,9 +13,12 @@ import ru.spcex.platform.imdg.iml.hazelcast.util.HazelcastHelper;
|
|||
|
||||
import java.util.List;
|
||||
import java.util.Random;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
|
||||
@Configuration
|
||||
public class HazelcastServiceTestConfiguration {
|
||||
|
||||
public static final AtomicLong currentID = new AtomicLong(0L);
|
||||
private HazelcastInstance hazelcastInstance;
|
||||
|
||||
private static ThreadPoolTaskExecutor createThreadPoolTaskExecutor(int maxPoolSz, boolean waitForCompletion) {
|
||||
|
|
|
|||
|
|
@ -3,8 +3,6 @@ package ru.spcex.clearing.company.service;
|
|||
import com.fasterxml.jackson.core.JsonProcessingException;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import com.hazelcast.core.IMap;
|
||||
import com.hazelcast.map.listener.EntryRemovedListener;
|
||||
import com.hazelcast.map.listener.EntryUpdatedListener;
|
||||
import org.apache.kafka.clients.consumer.ConsumerRecord;
|
||||
import org.apache.kafka.clients.consumer.MockConsumer;
|
||||
import org.apache.kafka.clients.consumer.OffsetResetStrategy;
|
||||
|
|
@ -19,20 +17,22 @@ import org.springframework.beans.factory.annotation.Qualifier;
|
|||
import org.springframework.test.context.ContextConfiguration;
|
||||
import org.springframework.test.context.junit.jupiter.SpringExtension;
|
||||
import ru.clearing.classes.statics.data.generated.ClearingMemberCategory;
|
||||
import ru.clearing.classes.statics.data.profile.Contact;
|
||||
import ru.spcex.clearing.company.config.HazelcastServiceTestConfiguration;
|
||||
import ru.spcex.clearing.company.utils.ImapEvent;
|
||||
import ru.spcex.clearing.company.utils.MatcherFactory.Matcher;
|
||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||
import ru.spcex.clearing.platform.messaging.domain.ActionType;
|
||||
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
|
||||
import ru.spcex.clearing.platform.messaging.domain.Consts;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteRequest;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.company.ClearingMemberCategoryNewRequest;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.company.ClearingMemberCategoryUpdateRequest;
|
||||
import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService;
|
||||
|
||||
import java.util.Collections;
|
||||
import java.util.HashMap;
|
||||
|
||||
import static ru.spcex.clearing.company.config.HazelcastServiceTestConfiguration.currentID;
|
||||
import static ru.spcex.clearing.company.utils.MatcherFactory.usingIgnoringFieldsComparator;
|
||||
|
||||
@ExtendWith(SpringExtension.class)
|
||||
|
|
@ -42,9 +42,10 @@ class ClearingMemberCategoryServiceTest {
|
|||
|
||||
public static final Matcher<ClearingMemberCategory> MEMBER_CATEGORY_MATCHER = usingIgnoringFieldsComparator();
|
||||
private static final int PARTITION = 0;
|
||||
private static final String TOPIC_MEMBER_CATEGORY_NEW = Consts.DESTINATION_CLEARING_MEMBER_CATEGORY_NEW;
|
||||
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 Long ID = 0L;
|
||||
private static final Long ID = currentID.getAndIncrement();
|
||||
|
||||
@Autowired
|
||||
@Qualifier("hazelcastServiceTest")
|
||||
|
|
@ -58,6 +59,56 @@ class ClearingMemberCategoryServiceTest {
|
|||
mockProducer = new MockProducer<>();
|
||||
}
|
||||
|
||||
@Test
|
||||
void clearingMemberCategoryNew() throws InterruptedException {
|
||||
//ARRANGE
|
||||
|
||||
ClearingMemberCategoryNewRequest memberCategoryNewRequest = new ClearingMemberCategoryNewRequest();
|
||||
memberCategoryNewRequest.setClearingMemberCategory("1234");
|
||||
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();
|
||||
predictableClearingMemberCategory.setId(ID);
|
||||
predictableClearingMemberCategory.setClearingMemberCategory("1234");
|
||||
|
||||
//ACT
|
||||
//service set up
|
||||
ClearingMemberCategoryService clearingMemberCategoryService = new ClearingMemberCategoryService(mockConsumer, mockProducer, hazelcastServiceTest);
|
||||
|
||||
//callbacks set up
|
||||
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
|
||||
ClearingMemberCategory resultUpdating = iMap.get(ID);
|
||||
MEMBER_CATEGORY_MATCHER.assertMatch(resultUpdating, predictableClearingMemberCategory);
|
||||
}
|
||||
|
||||
/**
|
||||
* {@link ClearingMemberCategoryService#clearingMemberCategoryUpdate(BaseRequest)}<br>
|
||||
* Тест проверяет обновление сущности {@link ClearingMemberCategory} в Hazelcast при передаче из Apache Kafka.<br>
|
||||
|
|
@ -100,6 +151,7 @@ class ClearingMemberCategoryServiceTest {
|
|||
|
||||
IMap<Long, ClearingMemberCategory> iMap = hazelcastServiceTest.getHazelcast().getMap(IMDGDistributedNames.Map_ClearingMemberCategory);
|
||||
iMap.put(ID, existsСlearingMemberCategory);
|
||||
ImapEvent imapEvent = new ImapEvent(iMap);
|
||||
|
||||
//KAFKA
|
||||
mockConsumer.schedulePollTask(() -> {
|
||||
|
|
@ -112,30 +164,11 @@ class ClearingMemberCategoryServiceTest {
|
|||
mockConsumer.updateBeginningOffsets(startOffsetsUpdating);
|
||||
|
||||
//waiting for hazelcast map item updates
|
||||
Object waiter = new Object();
|
||||
String listenerID = iMap.addEntryListener((EntryUpdatedListener<Long, Contact>) entryEvent -> {
|
||||
System.out.println("Checking If removed..");
|
||||
|
||||
synchronized (waiter) {
|
||||
try {
|
||||
waiter.wait(100);
|
||||
} catch (InterruptedException e) {
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
waiter.notify();
|
||||
}
|
||||
}, false);
|
||||
|
||||
synchronized (waiter) {
|
||||
waiter.wait(100);
|
||||
}
|
||||
imapEvent.waitWhenHappened();
|
||||
|
||||
//ASSERT
|
||||
ClearingMemberCategory resultUpdating = iMap.get(ID);
|
||||
MEMBER_CATEGORY_MATCHER.assertMatch(resultUpdating, predictableClearingMemberCategory);
|
||||
|
||||
//preparing hazelcastImdgProvider for next test
|
||||
iMap.removeEntryListener(listenerID);
|
||||
}
|
||||
|
||||
/**
|
||||
|
|
@ -175,6 +208,7 @@ class ClearingMemberCategoryServiceTest {
|
|||
|
||||
IMap<Long, ClearingMemberCategory> iMap = hazelcastServiceTest.getHazelcast().getMap(IMDGDistributedNames.Map_ClearingMemberCategory);
|
||||
iMap.put(ID, existsСlearingMemberCategory);
|
||||
ImapEvent imapEvent = new ImapEvent(iMap);
|
||||
|
||||
//KAFKA
|
||||
mockConsumer.schedulePollTask(() -> {
|
||||
|
|
@ -187,28 +221,9 @@ class ClearingMemberCategoryServiceTest {
|
|||
mockConsumer.updateBeginningOffsets(startOffsetsUpdating);
|
||||
|
||||
//waiting for hazelcast map item removes
|
||||
Object waiter = new Object();
|
||||
String listenerID = iMap.addEntryListener((EntryRemovedListener<Long, Contact>) entryEvent -> {
|
||||
System.out.println("Checking If removed..");
|
||||
|
||||
try {
|
||||
waiter.wait(100);
|
||||
} catch (InterruptedException e) {
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
synchronized (waiter) {
|
||||
waiter.notify();
|
||||
}
|
||||
}, false);
|
||||
|
||||
synchronized (waiter) {
|
||||
waiter.wait(100);
|
||||
}
|
||||
imapEvent.waitWhenHappened();
|
||||
|
||||
//ASSERT
|
||||
Assertions.assertEquals(0, iMap.size());
|
||||
|
||||
//preparing hazelcastImdgProvider for next test
|
||||
iMap.removeEntryListener(listenerID);
|
||||
}
|
||||
}
|
||||
|
|
@ -3,7 +3,6 @@ package ru.spcex.clearing.company.service;
|
|||
import com.fasterxml.jackson.core.JsonProcessingException;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import com.hazelcast.core.IMap;
|
||||
import com.hazelcast.map.listener.EntryUpdatedListener;
|
||||
import org.apache.kafka.clients.consumer.ConsumerRecord;
|
||||
import org.apache.kafka.clients.consumer.MockConsumer;
|
||||
import org.apache.kafka.clients.consumer.OffsetResetStrategy;
|
||||
|
|
@ -16,9 +15,10 @@ import org.springframework.beans.factory.annotation.Autowired;
|
|||
import org.springframework.beans.factory.annotation.Qualifier;
|
||||
import org.springframework.test.context.ContextConfiguration;
|
||||
import org.springframework.test.context.junit.jupiter.SpringExtension;
|
||||
import ru.clearing.classes.statics.data.company.Company;
|
||||
import ru.clearing.classes.statics.data.profile.CompanyInfo;
|
||||
import ru.clearing.classes.statics.data.profile.Contact;
|
||||
import ru.spcex.clearing.company.config.HazelcastServiceTestConfiguration;
|
||||
import ru.spcex.clearing.company.utils.ImapEvent;
|
||||
import ru.spcex.clearing.company.utils.MatcherFactory;
|
||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||
import ru.spcex.clearing.platform.messaging.domain.ActionType;
|
||||
|
|
@ -88,6 +88,9 @@ class CompanyInfoServiceTest {
|
|||
existsCompanyInfo.setFullNameEng("exists fullNameEng");
|
||||
existsCompanyInfo.setShortName("exists shortName");
|
||||
existsCompanyInfo.setFullName("exists fullName");
|
||||
Company existsCompany = new Company();
|
||||
existsCompany.setId(ID);
|
||||
existsCompany.setProfile(existsCompanyInfo);
|
||||
|
||||
CompanyInfoUpdateRequest companyInfoUpdateRequest = new CompanyInfoUpdateRequest();
|
||||
companyInfoUpdateRequest.setId(ID);
|
||||
|
|
@ -137,8 +140,9 @@ class CompanyInfoServiceTest {
|
|||
//callbacks set up
|
||||
clearingMemberCategoryService.afterPropertiesSet();
|
||||
|
||||
IMap<Long, CompanyInfo> iMap = hazelcastServiceTest.getHazelcast().getMap(IMDGDistributedNames.Map_Company);
|
||||
iMap.put(ID, existsCompanyInfo);
|
||||
IMap<Long, Company> iMap = hazelcastServiceTest.getHazelcast().getMap(IMDGDistributedNames.Map_Company);
|
||||
iMap.put(ID, existsCompany);
|
||||
ImapEvent imapEvent = new ImapEvent(iMap);
|
||||
|
||||
//KAFKA
|
||||
mockConsumer.schedulePollTask(() -> {
|
||||
|
|
@ -151,29 +155,10 @@ class CompanyInfoServiceTest {
|
|||
mockConsumer.updateBeginningOffsets(startOffsetsUpdating);
|
||||
|
||||
//waiting for hazelcast map updates
|
||||
Object waiter = new Object();
|
||||
String listenerID = iMap.addEntryListener((EntryUpdatedListener<Long, Contact>) entryEvent -> {
|
||||
System.out.println("Checking If removed..");
|
||||
|
||||
try {
|
||||
waiter.wait(100);
|
||||
} catch (InterruptedException e) {
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
synchronized (waiter) {
|
||||
waiter.notify();
|
||||
}
|
||||
}, false);
|
||||
|
||||
synchronized (waiter) {
|
||||
waiter.wait(100);
|
||||
}
|
||||
imapEvent.waitWhenHappened();
|
||||
|
||||
//ASSERT
|
||||
CompanyInfo resultUpdating = iMap.get(ID);
|
||||
COMPANY_INFO_MATCHER.assertMatch(resultUpdating, predictableCompanyInfo);
|
||||
|
||||
//preparing hazelcastImdgProvider for next test
|
||||
iMap.removeEntryListener(listenerID);
|
||||
Company resultUpdating = iMap.get(ID);
|
||||
COMPANY_INFO_MATCHER.assertMatch(resultUpdating.getProfile(), predictableCompanyInfo);
|
||||
}
|
||||
}
|
||||
|
|
@ -3,7 +3,6 @@ package ru.spcex.clearing.company.service;
|
|||
import com.fasterxml.jackson.core.JsonProcessingException;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import com.hazelcast.core.IMap;
|
||||
import com.hazelcast.map.listener.EntryRemovedListener;
|
||||
import org.apache.kafka.clients.consumer.ConsumerRecord;
|
||||
import org.apache.kafka.clients.consumer.MockConsumer;
|
||||
import org.apache.kafka.clients.consumer.OffsetResetStrategy;
|
||||
|
|
@ -19,8 +18,8 @@ import org.springframework.test.context.ContextConfiguration;
|
|||
import org.springframework.test.context.junit.jupiter.SpringExtension;
|
||||
import ru.clearing.classes.statics.data.company.Company;
|
||||
import ru.clearing.classes.statics.data.generated.ClearingMemberCategory;
|
||||
import ru.clearing.classes.statics.data.profile.Contact;
|
||||
import ru.spcex.clearing.company.config.HazelcastServiceTestConfiguration;
|
||||
import ru.spcex.clearing.company.utils.ImapEvent;
|
||||
import ru.spcex.clearing.company.utils.MatcherFactory;
|
||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||
import ru.spcex.clearing.platform.messaging.domain.ActionType;
|
||||
|
|
@ -92,6 +91,7 @@ class CompanyServiceTest {
|
|||
|
||||
IMap<Long, Company> iMap = hazelcastServiceTest.getHazelcast().getMap(IMDGDistributedNames.Map_Company);
|
||||
iMap.put(ID, existsCompany);
|
||||
ImapEvent imapEvent = new ImapEvent(iMap);
|
||||
|
||||
//KAFKA
|
||||
mockConsumer.schedulePollTask(() -> {
|
||||
|
|
@ -104,28 +104,9 @@ class CompanyServiceTest {
|
|||
mockConsumer.updateBeginningOffsets(startOffsetsUpdating);
|
||||
|
||||
//waiting for hazelcast map item removes
|
||||
Object waiter = new Object();
|
||||
String listenerID = iMap.addEntryListener((EntryRemovedListener<Long, Contact>) entryEvent -> {
|
||||
System.out.println("Checking If removed..");
|
||||
|
||||
try {
|
||||
waiter.wait(100);
|
||||
} catch (InterruptedException e) {
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
synchronized (waiter) {
|
||||
waiter.notify();
|
||||
}
|
||||
}, false);
|
||||
|
||||
synchronized (waiter) {
|
||||
waiter.wait(100);
|
||||
}
|
||||
imapEvent.waitWhenHappened();
|
||||
|
||||
//ASSERT
|
||||
Assertions.assertEquals(0, iMap.size());
|
||||
|
||||
//preparing hazelcastImdgProvider for next test
|
||||
iMap.removeEntryListener(listenerID);
|
||||
}
|
||||
}
|
||||
|
|
@ -3,7 +3,6 @@ package ru.spcex.clearing.company.service;
|
|||
import com.fasterxml.jackson.core.JsonProcessingException;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import com.hazelcast.core.IMap;
|
||||
import com.hazelcast.map.listener.EntryUpdatedListener;
|
||||
import org.apache.kafka.clients.consumer.ConsumerRecord;
|
||||
import org.apache.kafka.clients.consumer.MockConsumer;
|
||||
import org.apache.kafka.clients.consumer.OffsetResetStrategy;
|
||||
|
|
@ -17,8 +16,8 @@ import org.springframework.beans.factory.annotation.Qualifier;
|
|||
import org.springframework.test.context.ContextConfiguration;
|
||||
import org.springframework.test.context.junit.jupiter.SpringExtension;
|
||||
import ru.clearing.classes.statics.data.company.CompanySymbols;
|
||||
import ru.clearing.classes.statics.data.profile.Contact;
|
||||
import ru.spcex.clearing.company.config.HazelcastServiceTestConfiguration;
|
||||
import ru.spcex.clearing.company.utils.ImapEvent;
|
||||
import ru.spcex.clearing.company.utils.MatcherFactory;
|
||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||
import ru.spcex.clearing.platform.messaging.domain.ActionType;
|
||||
|
|
@ -105,6 +104,7 @@ class CompanySymbolServiceTest {
|
|||
|
||||
IMap<Long, CompanySymbols> iMap = hazelcastServiceTest.getHazelcast().getMap(IMDGDistributedNames.Map_CompanySymbols);
|
||||
iMap.put(ID, existsCompanySymbols);
|
||||
ImapEvent imapEvent = new ImapEvent(iMap);
|
||||
|
||||
//KAFKA
|
||||
mockConsumer.schedulePollTask(() -> {
|
||||
|
|
@ -117,29 +117,10 @@ class CompanySymbolServiceTest {
|
|||
mockConsumer.updateBeginningOffsets(startOffsetsUpdating);
|
||||
|
||||
//waiting for hazelcast map item updates
|
||||
Object waiter = new Object();
|
||||
String listenerID = iMap.addEntryListener((EntryUpdatedListener<Long, Contact>) entryEvent -> {
|
||||
System.out.println("Checking If removed..");
|
||||
|
||||
try {
|
||||
waiter.wait(100);
|
||||
} catch (InterruptedException e) {
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
synchronized (waiter) {
|
||||
waiter.notify();
|
||||
}
|
||||
}, false);
|
||||
|
||||
synchronized (waiter) {
|
||||
waiter.wait(100);
|
||||
}
|
||||
imapEvent.waitWhenHappened();
|
||||
|
||||
//ASSERT
|
||||
CompanySymbols resultUpdating = iMap.get(ID);
|
||||
COMPANY_SYMBOL_MATCHER.assertMatch(resultUpdating, predictableCompanySymbols);
|
||||
|
||||
//preparing hazelcastImdgProvider for next test
|
||||
iMap.removeEntryListener(listenerID);
|
||||
}
|
||||
}
|
||||
|
|
@ -3,7 +3,6 @@ package ru.spcex.clearing.company.service;
|
|||
import com.fasterxml.jackson.core.JsonProcessingException;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import com.hazelcast.core.IMap;
|
||||
import com.hazelcast.map.listener.EntryUpdatedListener;
|
||||
import org.apache.kafka.clients.consumer.ConsumerRecord;
|
||||
import org.apache.kafka.clients.consumer.MockConsumer;
|
||||
import org.apache.kafka.clients.consumer.OffsetResetStrategy;
|
||||
|
|
@ -18,6 +17,7 @@ import org.springframework.test.context.ContextConfiguration;
|
|||
import org.springframework.test.context.junit.jupiter.SpringExtension;
|
||||
import ru.clearing.classes.statics.data.profile.Contact;
|
||||
import ru.spcex.clearing.company.config.HazelcastServiceTestConfiguration;
|
||||
import ru.spcex.clearing.company.utils.ImapEvent;
|
||||
import ru.spcex.clearing.company.utils.MatcherFactory;
|
||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||
import ru.spcex.clearing.platform.messaging.domain.ActionType;
|
||||
|
|
@ -103,6 +103,7 @@ class ContactServiceTest {
|
|||
//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(() -> {
|
||||
|
|
@ -115,29 +116,10 @@ class ContactServiceTest {
|
|||
mockConsumer.updateBeginningOffsets(startOffsetsUpdating);
|
||||
|
||||
//waiting for hazelcast map updates
|
||||
Object waiter = new Object();
|
||||
String listenerID = iMap.addEntryListener((EntryUpdatedListener<Long, Contact>) entryEvent -> {
|
||||
System.out.println("Checking If removed..");
|
||||
|
||||
try {
|
||||
waiter.wait(100);
|
||||
} catch (InterruptedException e) {
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
synchronized (waiter) {
|
||||
waiter.notify();
|
||||
}
|
||||
}, false);
|
||||
|
||||
synchronized (waiter) {
|
||||
waiter.wait(100);
|
||||
}
|
||||
imapEvent.waitWhenHappened();
|
||||
|
||||
//ASSERT
|
||||
Contact resultUpdating = iMap.get(ID);
|
||||
CONTACT_MATCHER.assertMatch(resultUpdating, predictableContact);
|
||||
|
||||
//preparing hazelcastImdgProvider for next test
|
||||
iMap.removeEntryListener(listenerID);
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,63 @@
|
|||
package ru.spcex.clearing.company.utils;
|
||||
|
||||
import com.hazelcast.core.IMap;
|
||||
import com.hazelcast.map.listener.EntryAddedListener;
|
||||
import com.hazelcast.map.listener.EntryRemovedListener;
|
||||
import com.hazelcast.map.listener.EntryUpdatedListener;
|
||||
|
||||
import java.util.Timer;
|
||||
import java.util.TimerTask;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
|
||||
public class ImapEvent<T> {
|
||||
private final IMap<Long, T> iMap;
|
||||
private final String listenerAdding;
|
||||
private final String listenerUpdating;
|
||||
private final String listenerRemoving;
|
||||
private final AtomicBoolean checkEventHappened = new AtomicBoolean(false);
|
||||
|
||||
public ImapEvent(IMap<Long, T> iMap) {
|
||||
this.iMap = iMap;
|
||||
listenerAdding = iMap.addEntryListener((EntryAddedListener<Long, T>) entryEvent -> {
|
||||
synchronized (checkEventHappened) {
|
||||
checkEventHappened.set(true);
|
||||
checkEventHappened.notify();
|
||||
}
|
||||
}, false);
|
||||
listenerUpdating = iMap.addEntryListener((EntryUpdatedListener<Long, T>) entryEvent -> {
|
||||
synchronized (checkEventHappened) {
|
||||
checkEventHappened.set(true);
|
||||
checkEventHappened.notify();
|
||||
}
|
||||
}, false);
|
||||
listenerRemoving = iMap.addEntryListener((EntryRemovedListener<Long, T>) entryEvent -> {
|
||||
synchronized (checkEventHappened) {
|
||||
checkEventHappened.set(true);
|
||||
checkEventHappened.notify();
|
||||
}
|
||||
}, false);
|
||||
}
|
||||
|
||||
public void waitWhenHappened() throws InterruptedException {
|
||||
//running timer task as daemon thread
|
||||
Timer timer = new Timer(true);
|
||||
timer.scheduleAtFixedRate(new TimerTask() {
|
||||
boolean secondRan;
|
||||
|
||||
@Override
|
||||
public void run() {
|
||||
checkEventHappened.set(secondRan);//если что-то пойдет не так не тормозить основной поток
|
||||
secondRan = true;
|
||||
}
|
||||
}, 0, 30 * 1000);
|
||||
synchronized (checkEventHappened) {
|
||||
while (!checkEventHappened.get()) {
|
||||
checkEventHappened.wait(100);
|
||||
}
|
||||
}
|
||||
//preparing hazelcastImdgProvider for next test
|
||||
iMap.removeEntryListener(listenerAdding);
|
||||
iMap.removeEntryListener(listenerUpdating);
|
||||
iMap.removeEntryListener(listenerRemoving);
|
||||
}
|
||||
}
|
||||
Loading…
Add table
Reference in a new issue