From 2ad1b73300f9589325e3cbcc965f86b357f25d3e Mon Sep 17 00:00:00 2001 From: aalehin Date: Mon, 26 Dec 2022 16:55:51 +0300 Subject: [PATCH] =?UTF-8?q?http://jira.mfd.msk:8088/browse/CLS-78=20---=20?= =?UTF-8?q?=D0=9F=D1=80=D0=B8=20=D0=B4=D0=BE=D0=B1=D0=B0=D0=B2=D0=BB=D0=B5?= =?UTF-8?q?=D0=BD=D0=B8=D0=B8=20=D0=B7=D0=B0=D0=BF=D0=B8=D1=81=D0=B8=20?= =?UTF-8?q?=D0=B2=20plannerTemplate=20=D0=B8=20clearingCalendar=20=D0=B2?= =?UTF-8?q?=20=D0=BF=D0=BE=D0=BB=D0=B5=20createdAt=20=D1=81=D0=BE=D0=BE?= =?UTF-8?q?=D1=82=D0=B2=D0=B5=D1=82=D1=81=D1=82=D0=B2=D1=83=D1=8E=D1=89?= =?UTF-8?q?=D0=B8=D1=85=20=D1=82=D0=B0=D0=B1=D0=BB=D0=B8=D1=86=20=D0=BD?= =?UTF-8?q?=D0=B5=20=D0=B7=D0=B0=D0=BF=D0=BE=D0=BB=D0=BD=D1=8F=D0=BB=D0=BE?= =?UTF-8?q?=D1=81=D1=8C=20=D0=B2=D1=80=D0=B5=D0=BC=D1=8F=20=D1=81=D0=BE?= =?UTF-8?q?=D0=B7=D0=B4=D0=B0=D0=BD=D0=B8=D1=8F=20=D0=B7=D0=B0=D0=BF=D0=B8?= =?UTF-8?q?=D1=81=D0=B8.?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../service/ClearingCalendarService.java | 18 ++- .../service/PlannerTemplateService.java | 6 +- .../HazelcastServiceTestConfiguration.java | 68 +++++++++ .../service/ClearingCalendarServiceTest.java | 138 +++++++++++++++++ .../service/PlannerTemplateServiceTest.java | 141 ++++++++++++++++++ .../scheduler/utils/MatcherFactory.java | 38 +++++ 6 files changed, 401 insertions(+), 8 deletions(-) create mode 100644 clearing-parent/scheduler-service/src/test/java/ru/specx/clearing/scheduler/config/HazelcastServiceTestConfiguration.java create mode 100644 clearing-parent/scheduler-service/src/test/java/ru/specx/clearing/scheduler/service/ClearingCalendarServiceTest.java create mode 100644 clearing-parent/scheduler-service/src/test/java/ru/specx/clearing/scheduler/service/PlannerTemplateServiceTest.java create mode 100644 clearing-parent/scheduler-service/src/test/java/ru/specx/clearing/scheduler/utils/MatcherFactory.java diff --git a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/ClearingCalendarService.java b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/ClearingCalendarService.java index 70bd24147..3a86684b7 100644 --- a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/ClearingCalendarService.java +++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/ClearingCalendarService.java @@ -18,6 +18,8 @@ import ru.spcex.clearing.platform.messaging.service.QueueConsumer; import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.ImdgProvider; +import java.time.Instant; + @Service public class ClearingCalendarService extends QueueConsumer implements InitializingBean { private final Logger log = LoggerFactory.getLogger(getClass()); @@ -33,21 +35,23 @@ public class ClearingCalendarService extends QueueConsumer implements Initializi @Override public void afterPropertiesSet() { callback(ClearingCalendarNewRequest.class) - .setConsumer(this::newTradingCalendar) + .setConsumer(this::newClearingCalendar) .forDestination(Consts.DESTINATION_CLEARING_CALENDAR_NEW, callbacks::put); callback(ClearingCalendarUpdateRequest.class) - .setConsumer(this::updateTradingCalendar) + .setConsumer(this::updateClearingCalendar) .forDestination(Consts.DESTINATION_CLEARING_CALENDAR_UPDATE, callbacks::put); callback(CommonDeleteRequest.class) - .setConsumer(this::deleteTradingCalendar) + .setConsumer(this::deleteClearingCalendar) .forDestination(Consts.DESTINATION_CLEARING_CALENDAR_DELETE, callbacks::put); init(); } - private void newTradingCalendar(BaseRequest userRequest) { + private void newClearingCalendar(BaseRequest userRequest) { ClearingCalendarNewRequest req = userRequest.getRequestPayload(); log.debug("ClearingCalendarNewRequest received"); ClearingCalendar clearingCalendar = new ClearingCalendar(); + Instant created = Instant.now(); + clearingCalendar.setCreated(created); clearingCalendar.setClearingDate(req.getClearingDate()); clearingCalendar.setCompanyId(req.getCompanyId()); clearingCalendar.setDayStatus(req.getDayStatus()); @@ -55,18 +59,20 @@ public class ClearingCalendarService extends QueueConsumer implements Initializi log.debug("successfully processed, new id {}", clearingCalendar.getId()); } - private void updateTradingCalendar(BaseRequest userRequest) { + private void updateClearingCalendar(BaseRequest userRequest) { ClearingCalendarUpdateRequest req = userRequest.getRequestPayload(); log.debug("ClearingCalendarUpdateRequest received"); + Instant updated = Instant.now(); ClearingCalendar clearingCalendar = clearingCalendarMap.getSingleObjectByID(req.getId()); clearingCalendar.setClearingDate(req.getClearingDate()); clearingCalendar.setCompanyId(req.getCompanyId()); clearingCalendar.setDayStatus(req.getDayStatus()); + clearingCalendar.setUpdated(updated); clearingCalendarMap.update(clearingCalendar); log.debug("successfully processed, new id {}", clearingCalendar.getId()); } - private void deleteTradingCalendar(BaseRequest userRequest) { + private void deleteClearingCalendar(BaseRequest userRequest) { CommonDeleteRequest req = userRequest.getRequestPayload(); log.debug("CommonDeleteRequest received id = {}", req.getId()); ClearingCalendar clearingCalendar = clearingCalendarMap.getSingleObjectByID(req.getId()); diff --git a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/PlannerTemplateService.java b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/PlannerTemplateService.java index 491e884dc..77a1466a0 100644 --- a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/PlannerTemplateService.java +++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/PlannerTemplateService.java @@ -50,7 +50,8 @@ public class PlannerTemplateService extends QueueConsumer implements Initializin PlannerTemplateNewRequest req = userRequest.getRequestPayload(); log.debug("PlannerTemplateNewRequest received"); PlannerTemplate plannerTemplate = new PlannerTemplate(); - plannerTemplate.setCreated(Instant.now()); + Instant created = Instant.now(); + plannerTemplate.setCreated(created); plannerTemplate.setTask(req.getTask()); plannerTemplate.setTaskTime(req.getTaskTime()); plannerTemplate.setTaskStatus(req.getTaskStatus()); @@ -64,7 +65,8 @@ public class PlannerTemplateService extends QueueConsumer implements Initializin PlannerTemplateUpdateRequest req = userRequest.getRequestPayload(); log.debug("PlannerTemplateUpdateRequest received"); PlannerTemplate plannerTemplate = plannerTemplateMap.getSingleObjectByID(req.getId()); - plannerTemplate.setUpdated(Instant.now()); + Instant updated = Instant.now(); + plannerTemplate.setUpdated(updated); plannerTemplate.setTask(req.getTask()); plannerTemplate.setTaskTime(req.getTaskTime()); plannerTemplate.setTaskStatus(req.getTaskStatus()); diff --git a/clearing-parent/scheduler-service/src/test/java/ru/specx/clearing/scheduler/config/HazelcastServiceTestConfiguration.java b/clearing-parent/scheduler-service/src/test/java/ru/specx/clearing/scheduler/config/HazelcastServiceTestConfiguration.java new file mode 100644 index 000000000..af55150d0 --- /dev/null +++ b/clearing-parent/scheduler-service/src/test/java/ru/specx/clearing/scheduler/config/HazelcastServiceTestConfiguration.java @@ -0,0 +1,68 @@ +package ru.specx.clearing.scheduler.config; + +import com.hazelcast.config.*; +import com.hazelcast.core.Hazelcast; +import com.hazelcast.core.HazelcastInstance; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; +import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams; +import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService; +import ru.spcex.platform.imdg.iml.hazelcast.util.HazelcastHelper; + +import java.util.List; +import java.util.Random; + +@Configuration +public class HazelcastServiceTestConfiguration { + private HazelcastInstance hazelcastInstance; + + private static ThreadPoolTaskExecutor createThreadPoolTaskExecutor(int maxPoolSz, boolean waitForCompletion) { + ThreadPoolTaskExecutor pool = new ThreadPoolTaskExecutor(); + if (maxPoolSz > 2) { + pool.setKeepAliveSeconds(60); + pool.setAllowCoreThreadTimeOut(true); + } + pool.setCorePoolSize(maxPoolSz); + pool.setWaitForTasksToCompleteOnShutdown(waitForCompletion); + return pool; + } + + @Bean(name = "hazelcastServiceTest") + public HazelcastService hazelcastService(@Qualifier("taskExecutorHazelcastClientInitializer") ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer, @Qualifier("taskExecutorIdGeneratorAwaiter") ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter, HazelcastClientParams params) { + Config cfg = new Config(); + cfg.setInstanceName("localhost"); + + NetworkConfig networkConfig = new NetworkConfig(); + JoinConfig joinConfig = new JoinConfig(); + joinConfig.setMulticastConfig(new MulticastConfig().setEnabled(false)); + joinConfig.setTcpIpConfig(new TcpIpConfig().setEnabled(true).setMembers(List.of("127.0.0.1"))); + networkConfig.setJoin(joinConfig); + cfg.setNetworkConfig(networkConfig); + hazelcastInstance = Hazelcast.newHazelcastInstance(cfg); + HazelcastHelper.otcSystem_setStorageState(true, hazelcastInstance); + return new HazelcastService(taskExecutorHazelcastClientInitializer, taskExecutorIdGeneratorAwaiter, params); + } + + @Bean(name = "taskExecutorHazelcastClientInitializer") + public ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer() { + return createThreadPoolTaskExecutor(1, true); + } + + @Bean(name = "taskExecutorIdGeneratorAwaiter") + public ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter() { + return createThreadPoolTaskExecutor(1, false); + } + + @Bean(name = "hazelcastClientParams") + public HazelcastClientParams getHazelcastClientParams() { + HazelcastClientParams params = new HazelcastClientParams(); + params.setLogin("dev"); + params.setPassword("dev-pass"); + params.setClusterMembers("127.0.0.1"); + params.setInstanceName("hzTestClient" + new Random().nextInt()); + params.setNearCacheConfig(new NearCacheConfig()); + return params; + } +} diff --git a/clearing-parent/scheduler-service/src/test/java/ru/specx/clearing/scheduler/service/ClearingCalendarServiceTest.java b/clearing-parent/scheduler-service/src/test/java/ru/specx/clearing/scheduler/service/ClearingCalendarServiceTest.java new file mode 100644 index 000000000..7307966f9 --- /dev/null +++ b/clearing-parent/scheduler-service/src/test/java/ru/specx/clearing/scheduler/service/ClearingCalendarServiceTest.java @@ -0,0 +1,138 @@ +package ru.specx.clearing.scheduler.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; +import org.apache.kafka.clients.producer.MockProducer; +import org.apache.kafka.common.TopicPartition; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +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.profile.Contact; +import ru.clearing.classes.statics.data.scheduler.ClearingCalendar; +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.schedule.ClearingCalendarNewRequest; +import ru.spcex.clearing.scheduler.service.ClearingCalendarService; +import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService; +import ru.specx.clearing.scheduler.config.HazelcastServiceTestConfiguration; + +import java.time.Instant; +import java.time.LocalDate; +import java.util.Collections; +import java.util.HashMap; + +@ExtendWith(SpringExtension.class) +@ContextConfiguration(classes = { + HazelcastServiceTestConfiguration.class}) +class ClearingCalendarServiceTest { + + private static final int PARTITION = 0; + private static final String TOPIC_CLEARING_CALENDAR_NEW = Consts.DESTINATION_CLEARING_CALENDAR_NEW; + private static final Long ID = 0L; + + @Autowired + @Qualifier("hazelcastServiceTest") + private HazelcastService hazelcastServiceTest; + private MockConsumer mockConsumer; + private MockProducer mockProducer; + + @BeforeEach + void setUp() { + mockConsumer = new MockConsumer<>(OffsetResetStrategy.EARLIEST); + mockProducer = new MockProducer<>(); + } + + /** + * {@link ClearingCalendarService#newTradingCalendar(BaseRequest)}(BaseRequest)}
+ * Тест проверяет создание сущности {@link ClearingCalendar} в Hazelcast при передаче из Apache Kafka.
+ * Входной запрос {@link ClearingCalendarNewRequest}:
+ */ + @Test + public void newClearingCalendarInQueue() throws InterruptedException { + //ARRANGE + + ClearingCalendar clearingCalendar = new ClearingCalendar(); + Instant created = Instant.now(); + clearingCalendar.setCreated(created); + clearingCalendar.setClearingDate(LocalDate.now()); + clearingCalendar.setCompanyId(0L); + clearingCalendar.setDayStatus("STATUS"); + + ClearingCalendarNewRequest clearingCalendarNewRequest = new ClearingCalendarNewRequest(); + clearingCalendarNewRequest.setClearingDate(clearingCalendar.getClearingDate()); + clearingCalendarNewRequest.setCompanyId(clearingCalendar.getCompanyId()); + clearingCalendarNewRequest.setDayStatus(clearingCalendar.getDayStatus()); + + BaseRequest baseNewRequest = new BaseRequest<>(); + baseNewRequest.setRequestPayload(clearingCalendarNewRequest); + baseNewRequest.setId(ID); + baseNewRequest.setActionType(ActionType.NEW); + String jsonBaseForRequest; + ObjectMapper objectMapper = new ObjectMapper(); + try { + jsonBaseForRequest = objectMapper.writeValueAsString(baseNewRequest); + } catch (JsonProcessingException e) { + throw new RuntimeException(e); + } + + //ACT + //service set up + ClearingCalendarService clearingCalendarService = new ClearingCalendarService(mockConsumer, mockProducer, hazelcastServiceTest); + + //callbacks set up + clearingCalendarService.afterPropertiesSet(); + IMap iMap = hazelcastServiceTest.getHazelcast().getMap(IMDGDistributedNames.Map_ClearingCalendar); + //KAFKA + HashMap startOffsetsUpdating = new HashMap<>(); + TopicPartition topic = new TopicPartition(TOPIC_CLEARING_CALENDAR_NEW, PARTITION); + startOffsetsUpdating.put(topic, 0L); + mockConsumer.updateBeginningOffsets(startOffsetsUpdating); + + mockConsumer.schedulePollTask(() -> { + mockConsumer.rebalance(Collections.singletonList(new TopicPartition(TOPIC_CLEARING_CALENDAR_NEW, PARTITION))); + mockConsumer.addRecord(new ConsumerRecord<>(TOPIC_CLEARING_CALENDAR_NEW, PARTITION, 0, "key", jsonBaseForRequest)); + }); + + //waiting for hazelcast map item removes + Object waiter = new Object(); + String listenerID = iMap.addEntryListener((EntryRemovedListener) entryEvent -> { + System.out.println("Checking If pushed.."); + + try { + waiter.wait(100); + } catch (InterruptedException e) { + throw new RuntimeException(e); + } + synchronized (waiter) { + waiter.notify(); + } + }, false); + + synchronized (waiter) { + waiter.wait(100); + } + ClearingCalendar clearingCalendarRes = iMap.get(iMap.keySet().stream().findFirst().get()); + //ASSERT + Assertions.assertEquals(1, iMap.size()); + + Assertions.assertNotNull(clearingCalendarRes.getCreated()); + Assertions.assertEquals(clearingCalendar.getClearingDate(), clearingCalendarRes.getClearingDate()); + Assertions.assertEquals(clearingCalendar.getCompanyId(), clearingCalendarRes.getCompanyId()); + Assertions.assertEquals(clearingCalendar.getDayStatus(), clearingCalendarRes.getDayStatus()); + + //preparing hazelcastImdgProvider for next test + iMap.removeEntryListener(listenerID); + } +} \ No newline at end of file diff --git a/clearing-parent/scheduler-service/src/test/java/ru/specx/clearing/scheduler/service/PlannerTemplateServiceTest.java b/clearing-parent/scheduler-service/src/test/java/ru/specx/clearing/scheduler/service/PlannerTemplateServiceTest.java new file mode 100644 index 000000000..ff5ac8109 --- /dev/null +++ b/clearing-parent/scheduler-service/src/test/java/ru/specx/clearing/scheduler/service/PlannerTemplateServiceTest.java @@ -0,0 +1,141 @@ +package ru.specx.clearing.scheduler.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; +import org.apache.kafka.clients.producer.MockProducer; +import org.apache.kafka.common.TopicPartition; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +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.profile.Contact; +import ru.clearing.classes.statics.data.scheduler.PlannerTemplate; +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.schedule.PlannerTemplateNewRequest; +import ru.spcex.clearing.scheduler.service.PlannerTemplateService; +import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService; +import ru.specx.clearing.scheduler.config.HazelcastServiceTestConfiguration; + +import java.time.Instant; +import java.time.LocalTime; +import java.util.Collections; +import java.util.HashMap; + +@ExtendWith(SpringExtension.class) +@ContextConfiguration(classes = { + HazelcastServiceTestConfiguration.class}) +class PlannerTemplateServiceTest { + + private static final int PARTITION = 0; + private static final String TOPIC_PLANNER_TEMPLATE_NEW = Consts.DESTINATION_PLANNER_TEMPLATE_NEW; + private static final Long ID = 0L; + + @Autowired + @Qualifier("hazelcastServiceTest") + private HazelcastService hazelcastServiceTest; + private MockConsumer mockConsumer; + private MockProducer mockProducer; + + @BeforeEach + void setUp() { + mockConsumer = new MockConsumer<>(OffsetResetStrategy.EARLIEST); + mockProducer = new MockProducer<>(); + } + + /** + * {@link PlannerTemplateService#newTimetable(BaseRequest)}(BaseRequest)}
+ * Тест проверяет создание сущности {@link ru.clearing.classes.statics.data.scheduler.PlannerTemplate} в Hazelcast при передаче из Apache Kafka.
+ * Входной запрос {@link PlannerTemplateNewRequest}:
+ */ + @Test + public void newPlannerTemplateInQueue() throws InterruptedException { + //ARRANGE + Instant created = Instant.now(); + PlannerTemplate plannerTemplate = new PlannerTemplate(); + plannerTemplate.setTask("TASK"); + plannerTemplate.setTaskTime(LocalTime.MIDNIGHT); + plannerTemplate.setTaskStatus("TASK_STATUS"); + plannerTemplate.setCompanyId(0L); + plannerTemplate.setSecurityId(10L); + + PlannerTemplateNewRequest plannerTemplateNewRequest = new PlannerTemplateNewRequest(); + plannerTemplateNewRequest.setTask(plannerTemplate.getTask()); + plannerTemplateNewRequest.setTaskTime(plannerTemplate.getTaskTime()); + plannerTemplateNewRequest.setTaskStatus(plannerTemplate.getTaskStatus()); + plannerTemplateNewRequest.setCompanyId(plannerTemplate.getCompanyId()); + plannerTemplateNewRequest.setSecurityId(plannerTemplate.getSecurityId()); + + BaseRequest plannerTemplateRequest = new BaseRequest<>(); + plannerTemplateRequest.setRequestPayload(plannerTemplateNewRequest); + plannerTemplateRequest.setId(ID); + plannerTemplateRequest.setActionType(ActionType.NEW); + String jsonBaseForDeleteRequest; + ObjectMapper objectMapper = new ObjectMapper(); + try { + jsonBaseForDeleteRequest = objectMapper.writeValueAsString(plannerTemplateRequest); + } catch (JsonProcessingException e) { + throw new RuntimeException(e); + } + + //ACT + //service set up + PlannerTemplateService plannerTemplateService = new PlannerTemplateService(mockConsumer, mockProducer, hazelcastServiceTest); + + //callbacks set up + plannerTemplateService.afterPropertiesSet(); + IMap iMap = hazelcastServiceTest.getHazelcast().getMap(IMDGDistributedNames.Map_PlannerTemplate); + //KAFKA + HashMap startOffsetsUpdating = new HashMap<>(); + TopicPartition topic = new TopicPartition(TOPIC_PLANNER_TEMPLATE_NEW, PARTITION); + startOffsetsUpdating.put(topic, 0L); + mockConsumer.updateBeginningOffsets(startOffsetsUpdating); + + mockConsumer.schedulePollTask(() -> { + mockConsumer.rebalance(Collections.singletonList(new TopicPartition(TOPIC_PLANNER_TEMPLATE_NEW, PARTITION))); + mockConsumer.addRecord(new ConsumerRecord<>(TOPIC_PLANNER_TEMPLATE_NEW, PARTITION, 0, "key", jsonBaseForDeleteRequest)); + }); + + //waiting for hazelcast map item removes + Object waiter = new Object(); + String listenerID = iMap.addEntryListener((EntryRemovedListener) entryEvent -> { + System.out.println("Checking If pushed.."); + + try { + waiter.wait(100); + } catch (InterruptedException e) { + throw new RuntimeException(e); + } + synchronized (waiter) { + waiter.notify(); + } + }, false); + + synchronized (waiter) { + waiter.wait(100); + } + PlannerTemplate plannerTemplateRes = iMap.get(iMap.keySet().stream().findFirst().get()); + //ASSERT + Assertions.assertEquals(1, iMap.size()); + + Assertions.assertNotNull(plannerTemplateRes.getCreated()); + Assertions.assertEquals(plannerTemplate.getTask(), plannerTemplateRes.getTask()); + Assertions.assertEquals(plannerTemplate.getTaskTime(), plannerTemplateRes.getTaskTime()); + Assertions.assertEquals(plannerTemplate.getTaskStatus(), plannerTemplateRes.getTaskStatus()); + Assertions.assertEquals(plannerTemplate.getCompanyId(), plannerTemplateRes.getCompanyId()); + Assertions.assertEquals(plannerTemplate.getSecurityId(), plannerTemplateRes.getSecurityId()); + //preparing hazelcastImdgProvider for next test + iMap.removeEntryListener(listenerID); + } +} \ No newline at end of file diff --git a/clearing-parent/scheduler-service/src/test/java/ru/specx/clearing/scheduler/utils/MatcherFactory.java b/clearing-parent/scheduler-service/src/test/java/ru/specx/clearing/scheduler/utils/MatcherFactory.java new file mode 100644 index 000000000..c214393ea --- /dev/null +++ b/clearing-parent/scheduler-service/src/test/java/ru/specx/clearing/scheduler/utils/MatcherFactory.java @@ -0,0 +1,38 @@ +package ru.specx.clearing.scheduler.utils; + +import java.util.Arrays; + +import static org.assertj.core.api.Assertions.assertThat; + +/** + * Factory for creating test matchers. + *

+ * Comparing actual and expected objects via AssertJ + */ +public class MatcherFactory { + + public static Matcher usingIgnoringFieldsComparator(String... fieldsToIgnore) { + return new Matcher<>(fieldsToIgnore); + } + + public static class Matcher { + private final String[] fieldsToIgnore; + + private Matcher(String... fieldsToIgnore) { + this.fieldsToIgnore = fieldsToIgnore; + } + + public void assertMatch(T actual, T expected) { + assertThat(actual).usingRecursiveComparison().ignoringFields(fieldsToIgnore).isEqualTo(expected); + } + + @SafeVarargs + public final void assertMatch(Iterable actual, T... expected) { + assertMatch(actual, Arrays.asList(expected)); + } + + public void assertMatch(Iterable actual, Iterable expected) { + assertThat(actual).usingRecursiveFieldByFieldElementComparatorIgnoringFields(fieldsToIgnore).isEqualTo(expected); + } + } +}