--- При добавлении записи в plannerTemplate и clearingCalendar в поле createdAt соответствующих таблиц не заполнялось время создания записи.
This commit is contained in:
parent
a5320f36f3
commit
2ad1b73300
6 changed files with 401 additions and 8 deletions
|
|
@ -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.Imdg;
|
||||||
import ru.spcex.platform.imdg.api.ImdgProvider;
|
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||||
|
|
||||||
|
import java.time.Instant;
|
||||||
|
|
||||||
@Service
|
@Service
|
||||||
public class ClearingCalendarService extends QueueConsumer implements InitializingBean {
|
public class ClearingCalendarService extends QueueConsumer implements InitializingBean {
|
||||||
private final Logger log = LoggerFactory.getLogger(getClass());
|
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||||
|
|
@ -33,21 +35,23 @@ public class ClearingCalendarService extends QueueConsumer implements Initializi
|
||||||
@Override
|
@Override
|
||||||
public void afterPropertiesSet() {
|
public void afterPropertiesSet() {
|
||||||
callback(ClearingCalendarNewRequest.class)
|
callback(ClearingCalendarNewRequest.class)
|
||||||
.setConsumer(this::newTradingCalendar)
|
.setConsumer(this::newClearingCalendar)
|
||||||
.forDestination(Consts.DESTINATION_CLEARING_CALENDAR_NEW, callbacks::put);
|
.forDestination(Consts.DESTINATION_CLEARING_CALENDAR_NEW, callbacks::put);
|
||||||
callback(ClearingCalendarUpdateRequest.class)
|
callback(ClearingCalendarUpdateRequest.class)
|
||||||
.setConsumer(this::updateTradingCalendar)
|
.setConsumer(this::updateClearingCalendar)
|
||||||
.forDestination(Consts.DESTINATION_CLEARING_CALENDAR_UPDATE, callbacks::put);
|
.forDestination(Consts.DESTINATION_CLEARING_CALENDAR_UPDATE, callbacks::put);
|
||||||
callback(CommonDeleteRequest.class)
|
callback(CommonDeleteRequest.class)
|
||||||
.setConsumer(this::deleteTradingCalendar)
|
.setConsumer(this::deleteClearingCalendar)
|
||||||
.forDestination(Consts.DESTINATION_CLEARING_CALENDAR_DELETE, callbacks::put);
|
.forDestination(Consts.DESTINATION_CLEARING_CALENDAR_DELETE, callbacks::put);
|
||||||
init();
|
init();
|
||||||
}
|
}
|
||||||
|
|
||||||
private void newTradingCalendar(BaseRequest<ClearingCalendarNewRequest> userRequest) {
|
private void newClearingCalendar(BaseRequest<ClearingCalendarNewRequest> userRequest) {
|
||||||
ClearingCalendarNewRequest req = userRequest.getRequestPayload();
|
ClearingCalendarNewRequest req = userRequest.getRequestPayload();
|
||||||
log.debug("ClearingCalendarNewRequest received");
|
log.debug("ClearingCalendarNewRequest received");
|
||||||
ClearingCalendar clearingCalendar = new ClearingCalendar();
|
ClearingCalendar clearingCalendar = new ClearingCalendar();
|
||||||
|
Instant created = Instant.now();
|
||||||
|
clearingCalendar.setCreated(created);
|
||||||
clearingCalendar.setClearingDate(req.getClearingDate());
|
clearingCalendar.setClearingDate(req.getClearingDate());
|
||||||
clearingCalendar.setCompanyId(req.getCompanyId());
|
clearingCalendar.setCompanyId(req.getCompanyId());
|
||||||
clearingCalendar.setDayStatus(req.getDayStatus());
|
clearingCalendar.setDayStatus(req.getDayStatus());
|
||||||
|
|
@ -55,18 +59,20 @@ public class ClearingCalendarService extends QueueConsumer implements Initializi
|
||||||
log.debug("successfully processed, new id {}", clearingCalendar.getId());
|
log.debug("successfully processed, new id {}", clearingCalendar.getId());
|
||||||
}
|
}
|
||||||
|
|
||||||
private void updateTradingCalendar(BaseRequest<ClearingCalendarUpdateRequest> userRequest) {
|
private void updateClearingCalendar(BaseRequest<ClearingCalendarUpdateRequest> userRequest) {
|
||||||
ClearingCalendarUpdateRequest req = userRequest.getRequestPayload();
|
ClearingCalendarUpdateRequest req = userRequest.getRequestPayload();
|
||||||
log.debug("ClearingCalendarUpdateRequest received");
|
log.debug("ClearingCalendarUpdateRequest received");
|
||||||
|
Instant updated = Instant.now();
|
||||||
ClearingCalendar clearingCalendar = clearingCalendarMap.getSingleObjectByID(req.getId());
|
ClearingCalendar clearingCalendar = clearingCalendarMap.getSingleObjectByID(req.getId());
|
||||||
clearingCalendar.setClearingDate(req.getClearingDate());
|
clearingCalendar.setClearingDate(req.getClearingDate());
|
||||||
clearingCalendar.setCompanyId(req.getCompanyId());
|
clearingCalendar.setCompanyId(req.getCompanyId());
|
||||||
clearingCalendar.setDayStatus(req.getDayStatus());
|
clearingCalendar.setDayStatus(req.getDayStatus());
|
||||||
|
clearingCalendar.setUpdated(updated);
|
||||||
clearingCalendarMap.update(clearingCalendar);
|
clearingCalendarMap.update(clearingCalendar);
|
||||||
log.debug("successfully processed, new id {}", clearingCalendar.getId());
|
log.debug("successfully processed, new id {}", clearingCalendar.getId());
|
||||||
}
|
}
|
||||||
|
|
||||||
private void deleteTradingCalendar(BaseRequest<CommonDeleteRequest> userRequest) {
|
private void deleteClearingCalendar(BaseRequest<CommonDeleteRequest> userRequest) {
|
||||||
CommonDeleteRequest req = userRequest.getRequestPayload();
|
CommonDeleteRequest req = userRequest.getRequestPayload();
|
||||||
log.debug("CommonDeleteRequest received id = {}", req.getId());
|
log.debug("CommonDeleteRequest received id = {}", req.getId());
|
||||||
ClearingCalendar clearingCalendar = clearingCalendarMap.getSingleObjectByID(req.getId());
|
ClearingCalendar clearingCalendar = clearingCalendarMap.getSingleObjectByID(req.getId());
|
||||||
|
|
|
||||||
|
|
@ -50,7 +50,8 @@ public class PlannerTemplateService extends QueueConsumer implements Initializin
|
||||||
PlannerTemplateNewRequest req = userRequest.getRequestPayload();
|
PlannerTemplateNewRequest req = userRequest.getRequestPayload();
|
||||||
log.debug("PlannerTemplateNewRequest received");
|
log.debug("PlannerTemplateNewRequest received");
|
||||||
PlannerTemplate plannerTemplate = new PlannerTemplate();
|
PlannerTemplate plannerTemplate = new PlannerTemplate();
|
||||||
plannerTemplate.setCreated(Instant.now());
|
Instant created = Instant.now();
|
||||||
|
plannerTemplate.setCreated(created);
|
||||||
plannerTemplate.setTask(req.getTask());
|
plannerTemplate.setTask(req.getTask());
|
||||||
plannerTemplate.setTaskTime(req.getTaskTime());
|
plannerTemplate.setTaskTime(req.getTaskTime());
|
||||||
plannerTemplate.setTaskStatus(req.getTaskStatus());
|
plannerTemplate.setTaskStatus(req.getTaskStatus());
|
||||||
|
|
@ -64,7 +65,8 @@ public class PlannerTemplateService extends QueueConsumer implements Initializin
|
||||||
PlannerTemplateUpdateRequest req = userRequest.getRequestPayload();
|
PlannerTemplateUpdateRequest req = userRequest.getRequestPayload();
|
||||||
log.debug("PlannerTemplateUpdateRequest received");
|
log.debug("PlannerTemplateUpdateRequest received");
|
||||||
PlannerTemplate plannerTemplate = plannerTemplateMap.getSingleObjectByID(req.getId());
|
PlannerTemplate plannerTemplate = plannerTemplateMap.getSingleObjectByID(req.getId());
|
||||||
plannerTemplate.setUpdated(Instant.now());
|
Instant updated = Instant.now();
|
||||||
|
plannerTemplate.setUpdated(updated);
|
||||||
plannerTemplate.setTask(req.getTask());
|
plannerTemplate.setTask(req.getTask());
|
||||||
plannerTemplate.setTaskTime(req.getTaskTime());
|
plannerTemplate.setTaskTime(req.getTaskTime());
|
||||||
plannerTemplate.setTaskStatus(req.getTaskStatus());
|
plannerTemplate.setTaskStatus(req.getTaskStatus());
|
||||||
|
|
|
||||||
|
|
@ -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;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -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<String, Object> mockConsumer;
|
||||||
|
private MockProducer<String, Object> mockProducer;
|
||||||
|
|
||||||
|
@BeforeEach
|
||||||
|
void setUp() {
|
||||||
|
mockConsumer = new MockConsumer<>(OffsetResetStrategy.EARLIEST);
|
||||||
|
mockProducer = new MockProducer<>();
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* {@link ClearingCalendarService#newTradingCalendar(BaseRequest)}(BaseRequest)}<br>
|
||||||
|
* Тест проверяет создание сущности {@link ClearingCalendar} в Hazelcast при передаче из Apache Kafka.<br>
|
||||||
|
* Входной запрос {@link ClearingCalendarNewRequest}:<br>
|
||||||
|
*/
|
||||||
|
@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<ClearingCalendarNewRequest> 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<Long, ClearingCalendar> iMap = hazelcastServiceTest.getHazelcast().getMap(IMDGDistributedNames.Map_ClearingCalendar);
|
||||||
|
//KAFKA
|
||||||
|
HashMap<TopicPartition, Long> 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<Long, Contact>) 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);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -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<String, Object> mockConsumer;
|
||||||
|
private MockProducer<String, Object> mockProducer;
|
||||||
|
|
||||||
|
@BeforeEach
|
||||||
|
void setUp() {
|
||||||
|
mockConsumer = new MockConsumer<>(OffsetResetStrategy.EARLIEST);
|
||||||
|
mockProducer = new MockProducer<>();
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* {@link PlannerTemplateService#newTimetable(BaseRequest)}(BaseRequest)}<br>
|
||||||
|
* Тест проверяет создание сущности {@link ru.clearing.classes.statics.data.scheduler.PlannerTemplate} в Hazelcast при передаче из Apache Kafka.<br>
|
||||||
|
* Входной запрос {@link PlannerTemplateNewRequest}:<br>
|
||||||
|
*/
|
||||||
|
@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<PlannerTemplateNewRequest> 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<Long, PlannerTemplate> iMap = hazelcastServiceTest.getHazelcast().getMap(IMDGDistributedNames.Map_PlannerTemplate);
|
||||||
|
//KAFKA
|
||||||
|
HashMap<TopicPartition, Long> 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<Long, Contact>) 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);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -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.
|
||||||
|
* <p>
|
||||||
|
* Comparing actual and expected objects via AssertJ
|
||||||
|
*/
|
||||||
|
public class MatcherFactory {
|
||||||
|
|
||||||
|
public static <T> Matcher<T> usingIgnoringFieldsComparator(String... fieldsToIgnore) {
|
||||||
|
return new Matcher<>(fieldsToIgnore);
|
||||||
|
}
|
||||||
|
|
||||||
|
public static class Matcher<T> {
|
||||||
|
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<T> actual, T... expected) {
|
||||||
|
assertMatch(actual, Arrays.asList(expected));
|
||||||
|
}
|
||||||
|
|
||||||
|
public void assertMatch(Iterable<T> actual, Iterable<T> expected) {
|
||||||
|
assertThat(actual).usingRecursiveFieldByFieldElementComparatorIgnoringFields(fieldsToIgnore).isEqualTo(expected);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
Loading…
Add table
Reference in a new issue