Сделал TaskManager без hazelcast listener добавил отправку сообщений из сервисов в очередь и написал TaskManagerTest.
This commit is contained in:
parent
1d9d670c3a
commit
f41aecf739
7 changed files with 293 additions and 126 deletions
|
|
@ -0,0 +1,29 @@
|
||||||
|
package ru.spcex.clearing.scheduler.config;
|
||||||
|
|
||||||
|
import org.springframework.context.annotation.Bean;
|
||||||
|
import org.springframework.context.annotation.Configuration;
|
||||||
|
import ru.clearing.classes.statics.data.scheduler.PlannerAllToday;
|
||||||
|
import ru.spcex.clearing.scheduler.service.TaskManager;
|
||||||
|
|
||||||
|
import java.util.Map;
|
||||||
|
import java.util.concurrent.BlockingQueue;
|
||||||
|
import java.util.concurrent.LinkedBlockingQueue;
|
||||||
|
|
||||||
|
@Configuration
|
||||||
|
public class PlannerQueueConfig {
|
||||||
|
|
||||||
|
private static BlockingQueue<Map.Entry<TaskManager.Process, PlannerAllToday>> plannerQueue;
|
||||||
|
|
||||||
|
@Bean(name = "plannerQueue")
|
||||||
|
public BlockingQueue<Map.Entry<TaskManager.Process, PlannerAllToday>> plannerQueue() {
|
||||||
|
plannerQueue = new LinkedBlockingQueue<>();
|
||||||
|
return plannerQueue;
|
||||||
|
}
|
||||||
|
|
||||||
|
public static void addToPlannerQueue(TaskManager.Process process, PlannerAllToday plannerAllToday) {
|
||||||
|
try {
|
||||||
|
plannerQueue.put(Map.entry(process, plannerAllToday));
|
||||||
|
} catch (InterruptedException ignored) {
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -9,7 +9,6 @@ import org.springframework.beans.factory.annotation.Autowired;
|
||||||
import org.springframework.beans.factory.annotation.Qualifier;
|
import org.springframework.beans.factory.annotation.Qualifier;
|
||||||
import org.springframework.stereotype.Service;
|
import org.springframework.stereotype.Service;
|
||||||
import ru.clearing.classes.statics.data.scheduler.ClearingCalendar;
|
import ru.clearing.classes.statics.data.scheduler.ClearingCalendar;
|
||||||
import ru.clearing.classes.statics.data.scheduler.Planner;
|
|
||||||
import ru.clearing.classes.statics.data.scheduler.PlannerAllToday;
|
import ru.clearing.classes.statics.data.scheduler.PlannerAllToday;
|
||||||
import ru.clearing.classes.statics.data.scheduler.PlannerTemplate;
|
import ru.clearing.classes.statics.data.scheduler.PlannerTemplate;
|
||||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||||
|
|
@ -22,9 +21,8 @@ import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
|
||||||
import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate;
|
import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate;
|
||||||
import ru.spcex.clearing.scheduler.PlannerAllTodayBuilder;
|
import ru.spcex.clearing.scheduler.PlannerAllTodayBuilder;
|
||||||
import ru.spcex.clearing.util.security.UserRoleVerification;
|
import ru.spcex.clearing.util.security.UserRoleVerification;
|
||||||
import ru.spcex.platform.classes.base.SpcexObjectBase;
|
|
||||||
import ru.spcex.platform.enumeration.DayStatus;
|
import ru.spcex.platform.enumeration.DayStatus;
|
||||||
import ru.spcex.platform.enumeration.Status;
|
import ru.spcex.platform.enumeration.Parent;
|
||||||
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 ru.spcex.platform.utils.enumeration.IMessageResolver;
|
import ru.spcex.platform.utils.enumeration.IMessageResolver;
|
||||||
|
|
@ -35,11 +33,11 @@ import java.time.Instant;
|
||||||
import java.time.LocalDate;
|
import java.time.LocalDate;
|
||||||
import java.util.Arrays;
|
import java.util.Arrays;
|
||||||
import java.util.Collection;
|
import java.util.Collection;
|
||||||
import java.util.List;
|
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
import java.util.function.Function;
|
import java.util.function.Function;
|
||||||
|
|
||||||
import static ru.spcex.clearing.scheduler.ISchedulerChecker.isValidWorkday;
|
import static ru.spcex.clearing.scheduler.ISchedulerChecker.isValidWorkday;
|
||||||
|
import static ru.spcex.clearing.scheduler.config.PlannerQueueConfig.addToPlannerQueue;
|
||||||
|
|
||||||
@Service
|
@Service
|
||||||
public class ClearingCalendarService extends QueueConsumer implements InitializingBean {
|
public class ClearingCalendarService extends QueueConsumer implements InitializingBean {
|
||||||
|
|
@ -47,7 +45,6 @@ public class ClearingCalendarService extends QueueConsumer implements Initializi
|
||||||
private final Imdg<ClearingCalendar> clearingCalendarMap;
|
private final Imdg<ClearingCalendar> clearingCalendarMap;
|
||||||
private final Imdg<PlannerTemplate> plannerTemplateMap;
|
private final Imdg<PlannerTemplate> plannerTemplateMap;
|
||||||
private final Imdg<PlannerAllToday> plannerAllTodayMap;
|
private final Imdg<PlannerAllToday> plannerAllTodayMap;
|
||||||
private final Imdg<Planner> plannerMap;
|
|
||||||
private final IMessageResolver messageResolver;
|
private final IMessageResolver messageResolver;
|
||||||
private final UserRoleVerification userRoleVerification;
|
private final UserRoleVerification userRoleVerification;
|
||||||
private final Function<CommonDeleteRequest, IValidator> clearingCalendarDeleteRequestValidation;
|
private final Function<CommonDeleteRequest, IValidator> clearingCalendarDeleteRequestValidation;
|
||||||
|
|
@ -65,7 +62,6 @@ public class ClearingCalendarService extends QueueConsumer implements Initializi
|
||||||
super(kafkaQueue, kafkaProducer);
|
super(kafkaQueue, kafkaProducer);
|
||||||
this.clearingCalendarMap = imdgProvider.getImdg(IMDGDistributedNames.Map_ClearingCalendar, ClearingCalendar.class);
|
this.clearingCalendarMap = imdgProvider.getImdg(IMDGDistributedNames.Map_ClearingCalendar, ClearingCalendar.class);
|
||||||
this.plannerTemplateMap = imdgProvider.getImdg(IMDGDistributedNames.Map_PlannerTemplate, PlannerTemplate.class);
|
this.plannerTemplateMap = imdgProvider.getImdg(IMDGDistributedNames.Map_PlannerTemplate, PlannerTemplate.class);
|
||||||
this.plannerMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Planner, Planner.class);
|
|
||||||
this.plannerAllTodayMap = imdgProvider.getImdg(IMDGDistributedNames.Map_PlannerAllToday, PlannerAllToday.class);
|
this.plannerAllTodayMap = imdgProvider.getImdg(IMDGDistributedNames.Map_PlannerAllToday, PlannerAllToday.class);
|
||||||
this.userRoleVerification = userRoleVerification;
|
this.userRoleVerification = userRoleVerification;
|
||||||
this.clearingCalendarDeleteRequestValidation = clearingCalendarDeleteRequestValidator;
|
this.clearingCalendarDeleteRequestValidation = clearingCalendarDeleteRequestValidator;
|
||||||
|
|
@ -157,8 +153,11 @@ public class ClearingCalendarService extends QueueConsumer implements Initializi
|
||||||
if (isValidWorkday(clearingCalendar, isWeekend)) {
|
if (isValidWorkday(clearingCalendar, isWeekend)) {
|
||||||
for (PlannerTemplate plannerTemplate : plannerTemplateMap.getAllValues()) {
|
for (PlannerTemplate plannerTemplate : plannerTemplateMap.getAllValues()) {
|
||||||
PlannerAllToday plannerAllToday = plannerAllTodayMap.getSingleObjectBySQL(String.format("parentId = %s", plannerTemplate.getId()));
|
PlannerAllToday plannerAllToday = plannerAllTodayMap.getSingleObjectBySQL(String.format("parentId = %s", plannerTemplate.getId()));
|
||||||
if (plannerAllToday == null)
|
if (plannerAllToday == null) {
|
||||||
plannerAllTodayMap.insert(PlannerAllTodayBuilder.builder().append(plannerTemplate).build());
|
plannerAllToday = PlannerAllTodayBuilder.builder().append(plannerTemplate).build();
|
||||||
|
plannerAllTodayMap.insert(plannerAllToday);
|
||||||
|
addToPlannerQueue(TaskManager.Process.add, plannerAllToday);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -177,14 +176,11 @@ public class ClearingCalendarService extends QueueConsumer implements Initializi
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
//удалим все PlannerAllToday созданные на основании plannerTemplate кроме созданных на основании planner
|
//удалим все PlannerAllToday созданные на основании plannerTemplate кроме созданных на основании planner
|
||||||
List<Long> plannerIds = plannerMap.getCollectionObjectsByFieldValues(Map.of("taskStatus", Status.Active.getKey(),
|
Collection<PlannerAllToday> values = plannerAllTodayMap.getCollectionObjectsByFieldValues(Map.of("parent", Parent.Template.getKey()));
|
||||||
"clearingDate", currentDate)).stream()
|
|
||||||
.mapToLong(SpcexObjectBase::getId)
|
|
||||||
.boxed().toList();
|
|
||||||
|
|
||||||
Collection<PlannerAllToday> values = plannerAllTodayMap.getAllValues().stream()
|
values.forEach(plannerAllToday -> {
|
||||||
.filter(plannerAllToday -> !plannerIds.contains(plannerAllToday.getParentId())).toList();
|
plannerAllTodayMap.delete(plannerAllToday);
|
||||||
|
addToPlannerQueue(TaskManager.Process.delete, plannerAllToday);
|
||||||
values.forEach(plannerAllTodayMap::delete);
|
});
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -39,6 +39,7 @@ public class LauncherSender {
|
||||||
BaseRequest<Object> request = new BaseRequest<>();
|
BaseRequest<Object> request = new BaseRequest<>();
|
||||||
request.setId(idGenerator.nextId());
|
request.setId(idGenerator.nextId());
|
||||||
request.setActionType(ActionType.NEW);
|
request.setActionType(ActionType.NEW);
|
||||||
|
request.setUserId(userId);
|
||||||
request.setRequestPayload(toRequest);// iAction.toRequest()
|
request.setRequestPayload(toRequest);// iAction.toRequest()
|
||||||
return request;
|
return request;
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -32,6 +32,8 @@ import java.util.Collection;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
import java.util.function.Function;
|
import java.util.function.Function;
|
||||||
|
|
||||||
|
import static ru.spcex.clearing.scheduler.config.PlannerQueueConfig.addToPlannerQueue;
|
||||||
|
|
||||||
@Service
|
@Service
|
||||||
public class PlannerService extends QueueConsumer implements InitializingBean {
|
public class PlannerService extends QueueConsumer implements InitializingBean {
|
||||||
private final Logger log = LoggerFactory.getLogger(getClass());
|
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||||
|
|
@ -151,9 +153,16 @@ public class PlannerService extends QueueConsumer implements InitializingBean {
|
||||||
public void cudPlannerAllToday(Planner planner) {
|
public void cudPlannerAllToday(Planner planner) {
|
||||||
if (planner.getTaskStatus().equalsIgnoreCase(Status.Active.getKey())) {
|
if (planner.getTaskStatus().equalsIgnoreCase(Status.Active.getKey())) {
|
||||||
PlannerAllToday plannerAllToday = PlannerAllTodayBuilder.builder().append(planner).build();
|
PlannerAllToday plannerAllToday = PlannerAllTodayBuilder.builder().append(planner).build();
|
||||||
PlannerAllToday allToday = plannerAllTodayMap.getSingleObjectBySQL(String.format("parentId = %s", planner.getId()));
|
PlannerAllToday oldPlannerAllToday = plannerAllTodayMap.getSingleObjectBySQL(String.format("parentId = %s", planner.getId()));
|
||||||
if (allToday != null) plannerAllToday.setId(allToday.getId());
|
if (oldPlannerAllToday != null) {
|
||||||
plannerAllTodayMap.insert(plannerAllToday);
|
plannerAllToday.setId(oldPlannerAllToday.getId());
|
||||||
|
plannerAllTodayMap.insert(plannerAllToday);
|
||||||
|
addToPlannerQueue(TaskManager.Process.update, oldPlannerAllToday);
|
||||||
|
} else {
|
||||||
|
plannerAllTodayMap.insert(plannerAllToday);
|
||||||
|
addToPlannerQueue(TaskManager.Process.add, plannerAllToday);
|
||||||
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
||||||
if (planner.getTaskStatus().equalsIgnoreCase(Status.Cancel.getKey()) || planner.getTaskStatus().equalsIgnoreCase(Status.Blocked.getKey())) {
|
if (planner.getTaskStatus().equalsIgnoreCase(Status.Cancel.getKey()) || planner.getTaskStatus().equalsIgnoreCase(Status.Blocked.getKey())) {
|
||||||
|
|
@ -162,7 +171,10 @@ public class PlannerService extends QueueConsumer implements InitializingBean {
|
||||||
"taskTime", planner.getTaskTime(),
|
"taskTime", planner.getTaskTime(),
|
||||||
"securityId", planner.getSecurityId(),
|
"securityId", planner.getSecurityId(),
|
||||||
"companyId", planner.getCompanyId()));
|
"companyId", planner.getCompanyId()));
|
||||||
values.forEach(plannerAllTodayMap::delete);
|
values.forEach(plannerAllToday -> {
|
||||||
|
plannerAllTodayMap.delete(plannerAllToday);
|
||||||
|
addToPlannerQueue(TaskManager.Process.delete, plannerAllToday);
|
||||||
|
});
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -36,6 +36,7 @@ import java.util.Map;
|
||||||
import java.util.function.Function;
|
import java.util.function.Function;
|
||||||
|
|
||||||
import static ru.spcex.clearing.scheduler.ISchedulerChecker.isValidWorkday;
|
import static ru.spcex.clearing.scheduler.ISchedulerChecker.isValidWorkday;
|
||||||
|
import static ru.spcex.clearing.scheduler.config.PlannerQueueConfig.addToPlannerQueue;
|
||||||
|
|
||||||
@Service
|
@Service
|
||||||
public class PlannerTemplateService extends QueueConsumer implements InitializingBean {
|
public class PlannerTemplateService extends QueueConsumer implements InitializingBean {
|
||||||
|
|
@ -102,7 +103,11 @@ public class PlannerTemplateService extends QueueConsumer implements Initializin
|
||||||
plannerTemplate.setSecurityId(req.getSecurityId());
|
plannerTemplate.setSecurityId(req.getSecurityId());
|
||||||
plannerTemplateMap.insert(plannerTemplate);
|
plannerTemplateMap.insert(plannerTemplate);
|
||||||
|
|
||||||
if (todayWorkDay()) plannerAllTodayMap.insert(PlannerAllTodayBuilder.builder().append(plannerTemplate).build());
|
if (todayWorkDay()) {
|
||||||
|
PlannerAllToday plannerAllToday = PlannerAllTodayBuilder.builder().append(plannerTemplate).build();
|
||||||
|
plannerAllTodayMap.insert(plannerAllToday);
|
||||||
|
addToPlannerQueue(TaskManager.Process.add, plannerAllToday);
|
||||||
|
}
|
||||||
log.debug("successfully processed, new id {}", plannerTemplate.getId());
|
log.debug("successfully processed, new id {}", plannerTemplate.getId());
|
||||||
return null;
|
return null;
|
||||||
}
|
}
|
||||||
|
|
@ -127,10 +132,16 @@ public class PlannerTemplateService extends QueueConsumer implements Initializin
|
||||||
plannerTemplateMap.update(plannerTemplate);
|
plannerTemplateMap.update(plannerTemplate);
|
||||||
|
|
||||||
if (todayWorkDay()) {
|
if (todayWorkDay()) {
|
||||||
PlannerAllToday planner = PlannerAllTodayBuilder.builder().append(plannerTemplate).build();
|
PlannerAllToday plannerAllToday = PlannerAllTodayBuilder.builder().append(plannerTemplate).build();
|
||||||
PlannerAllToday plannerAllToday = plannerAllTodayMap.getSingleObjectBySQL(String.format("parentId = %s", plannerTemplate.getId()));
|
PlannerAllToday oldPlannerAllToday = plannerAllTodayMap.getSingleObjectBySQL(String.format("parentId = %s", plannerTemplate.getId()));
|
||||||
if (plannerAllToday != null) planner.setId(plannerAllToday.getId());
|
if (oldPlannerAllToday != null) {
|
||||||
plannerAllTodayMap.insert(planner);
|
plannerAllToday.setId(oldPlannerAllToday.getId());
|
||||||
|
plannerAllTodayMap.insert(plannerAllToday);
|
||||||
|
addToPlannerQueue(TaskManager.Process.update, oldPlannerAllToday);
|
||||||
|
}else {
|
||||||
|
plannerAllTodayMap.insert(plannerAllToday);
|
||||||
|
addToPlannerQueue(TaskManager.Process.add, plannerAllToday);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
log.debug("successfully processed, new id {}", plannerTemplate.getId());
|
log.debug("successfully processed, new id {}", plannerTemplate.getId());
|
||||||
return null;
|
return null;
|
||||||
|
|
@ -149,7 +160,10 @@ public class PlannerTemplateService extends QueueConsumer implements Initializin
|
||||||
|
|
||||||
if (todayWorkDay()) {
|
if (todayWorkDay()) {
|
||||||
PlannerAllToday plannerAllToday = plannerAllTodayMap.getSingleObjectBySQL(String.format("parentId = %s", plannerTemplate.getId()));
|
PlannerAllToday plannerAllToday = plannerAllTodayMap.getSingleObjectBySQL(String.format("parentId = %s", plannerTemplate.getId()));
|
||||||
if (plannerAllToday != null) plannerAllTodayMap.delete(plannerAllToday);
|
if (plannerAllToday != null){
|
||||||
|
plannerAllTodayMap.delete(plannerAllToday);
|
||||||
|
addToPlannerQueue(TaskManager.Process.delete, plannerAllToday);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
plannerTemplateMap.delete(plannerTemplate);
|
plannerTemplateMap.delete(plannerTemplate);
|
||||||
return null;
|
return null;
|
||||||
|
|
|
||||||
|
|
@ -1,14 +1,10 @@
|
||||||
package ru.spcex.clearing.scheduler.service;
|
package ru.spcex.clearing.scheduler.service;
|
||||||
|
|
||||||
import com.hazelcast.core.EntryEvent;
|
|
||||||
import com.hazelcast.map.listener.EntryAddedListener;
|
|
||||||
import com.hazelcast.map.listener.EntryRemovedListener;
|
|
||||||
import com.hazelcast.map.listener.EntryUpdatedListener;
|
|
||||||
import org.slf4j.Logger;
|
import org.slf4j.Logger;
|
||||||
import org.slf4j.LoggerFactory;
|
import org.slf4j.LoggerFactory;
|
||||||
import org.springframework.beans.factory.InitializingBean;
|
import org.springframework.beans.factory.InitializingBean;
|
||||||
import org.springframework.beans.factory.annotation.Autowired;
|
import org.springframework.beans.factory.annotation.Autowired;
|
||||||
import org.springframework.context.annotation.Lazy;
|
import org.springframework.beans.factory.annotation.Qualifier;
|
||||||
import org.springframework.lang.NonNull;
|
import org.springframework.lang.NonNull;
|
||||||
import org.springframework.scheduling.TaskScheduler;
|
import org.springframework.scheduling.TaskScheduler;
|
||||||
import org.springframework.stereotype.Service;
|
import org.springframework.stereotype.Service;
|
||||||
|
|
@ -18,15 +14,12 @@ import ru.spcex.platform.enumeration.Status;
|
||||||
import ru.spcex.platform.enumeration.Task;
|
import ru.spcex.platform.enumeration.Task;
|
||||||
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 ru.spcex.platform.imdg.iml.hazelcast.adapter.ImdgHazelcast;
|
|
||||||
|
|
||||||
import java.time.*;
|
import java.time.*;
|
||||||
import java.util.ArrayList;
|
import java.util.*;
|
||||||
import java.util.Collection;
|
import java.util.concurrent.*;
|
||||||
import java.util.Date;
|
import java.util.concurrent.atomic.AtomicBoolean;
|
||||||
import java.util.Objects;
|
import java.util.function.Consumer;
|
||||||
import java.util.concurrent.ConcurrentHashMap;
|
|
||||||
import java.util.concurrent.ScheduledFuture;
|
|
||||||
import java.util.stream.Collectors;
|
import java.util.stream.Collectors;
|
||||||
|
|
||||||
import static ru.spcex.clearing.imdg.IMDGDistributedNames.Map_Launcher;
|
import static ru.spcex.clearing.imdg.IMDGDistributedNames.Map_Launcher;
|
||||||
|
|
@ -39,25 +32,63 @@ import static ru.spcex.platform.utils.enumeration.IEnumKey.getEnumByKey;
|
||||||
* <p>
|
* <p>
|
||||||
*/
|
*/
|
||||||
@Service("userReportTask")
|
@Service("userReportTask")
|
||||||
@Lazy
|
public class TaskManager implements InitializingBean, AutoCloseable {
|
||||||
public class TaskManager implements EntryAddedListener<Long, PlannerAllToday>,
|
|
||||||
EntryUpdatedListener<Long, PlannerAllToday>, EntryRemovedListener<Long, PlannerAllToday>,
|
|
||||||
InitializingBean {
|
|
||||||
private static final Logger log = LoggerFactory.getLogger(TaskManager.class);
|
private static final Logger log = LoggerFactory.getLogger(TaskManager.class);
|
||||||
|
public static final Long systemId = 0L;//для системных задач по договоренности, должен быть пользователь с правами админа и id=0L
|
||||||
|
|
||||||
private final TaskScheduler taskScheduler;
|
private final TaskScheduler taskScheduler;
|
||||||
private final ImdgProvider imdgProvider;
|
private final ImdgProvider imdgProvider;
|
||||||
private Imdg<Launcher> launcherMap;
|
private Imdg<Launcher> launcherMap;
|
||||||
private Imdg<PlannerAllToday> plannerAllTodayMap;
|
private Imdg<PlannerAllToday> plannerAllTodayMap;
|
||||||
private ConcurrentHashMap<LocalTime, ScheduledFuture> scheduledJobs;
|
private ConcurrentHashMap<Long, ScheduledFuture> scheduledJobs;
|
||||||
private LauncherSender launcherSender;
|
private LauncherSender launcherSender;
|
||||||
|
private final BlockingQueue<Map.Entry<Process, PlannerAllToday>> plannerQueue;
|
||||||
|
private final Map<Process, Consumer<PlannerAllToday>> callbacks = new HashMap<>();
|
||||||
|
private final AtomicBoolean closed = new AtomicBoolean(false);
|
||||||
|
private final ExecutorService outputExecutor;
|
||||||
|
|
||||||
@Autowired
|
@Autowired
|
||||||
TaskManager(TaskScheduler taskScheduler,
|
TaskManager(TaskScheduler taskScheduler,
|
||||||
ImdgProvider imdgProvider,
|
ImdgProvider imdgProvider,
|
||||||
LauncherSender launcherSender) {
|
LauncherSender launcherSender,
|
||||||
|
@Qualifier("plannerQueue") BlockingQueue<Map.Entry<Process, PlannerAllToday>> plannerQueue) {
|
||||||
this.taskScheduler = taskScheduler;
|
this.taskScheduler = taskScheduler;
|
||||||
this.imdgProvider = imdgProvider;
|
this.imdgProvider = imdgProvider;
|
||||||
this.launcherSender = launcherSender;
|
this.launcherSender = launcherSender;
|
||||||
|
this.plannerQueue = plannerQueue;
|
||||||
|
this.outputExecutor = Executors.newSingleThreadExecutor();
|
||||||
|
}
|
||||||
|
|
||||||
|
public enum Process {
|
||||||
|
add, update, delete
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public void afterPropertiesSet() {
|
||||||
|
this.launcherMap = imdgProvider.getImdg(Map_Launcher, Launcher.class);
|
||||||
|
this.plannerAllTodayMap = imdgProvider.getImdg(Map_PlannerAllToday, PlannerAllToday.class);
|
||||||
|
callbacks.put(Process.add, this::entryAdded);
|
||||||
|
callbacks.put(Process.update, this::entryUpdated);
|
||||||
|
callbacks.put(Process.delete, this::entryRemoved);
|
||||||
|
scheduledJobs = new ConcurrentHashMap<>();
|
||||||
|
//если сервис стартовал раньше imdg, дождемся создания всех PlannerAllToday.
|
||||||
|
imdgProvider.waitAvailable();
|
||||||
|
updateScheduler();
|
||||||
|
process();
|
||||||
|
}
|
||||||
|
|
||||||
|
private void process() {
|
||||||
|
outputExecutor.submit(() -> {
|
||||||
|
try {
|
||||||
|
while (!closed.get()) {
|
||||||
|
//процесс ожидает пока появится новое сообщение в очереди
|
||||||
|
Map.Entry<Process, PlannerAllToday> entry = plannerQueue.take();
|
||||||
|
callbacks.get(entry.getKey()).accept(entry.getValue());
|
||||||
|
}
|
||||||
|
} catch (InterruptedException e) {
|
||||||
|
// closed.set(true);
|
||||||
|
}
|
||||||
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
private static LocalDateTime dateOldTypeConvert(@NonNull Date oldDate) {
|
private static LocalDateTime dateOldTypeConvert(@NonNull Date oldDate) {
|
||||||
|
|
@ -72,85 +103,65 @@ public class TaskManager implements EntryAddedListener<Long, PlannerAllToday>,
|
||||||
return dateOldTypeConvert(oldDate).toLocalTime();
|
return dateOldTypeConvert(oldDate).toLocalTime();
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
// --- Обработчики PlannerAllToday ---
|
||||||
public void afterPropertiesSet() {
|
public void entryAdded(PlannerAllToday task) {
|
||||||
this.launcherMap = imdgProvider.getImdg(Map_Launcher, Launcher.class);
|
String taskStatus = task.getTaskStatus();
|
||||||
this.plannerAllTodayMap = imdgProvider.getImdg(Map_PlannerAllToday, PlannerAllToday.class);
|
if (taskStatus == null) {
|
||||||
if (plannerAllTodayMap instanceof ImdgHazelcast<PlannerAllToday> plannerAllTodayImdgHazelcast) {
|
log.warn("Task skipped, {} task status not recognized", task.getTask());
|
||||||
plannerAllTodayImdgHazelcast.getMap().addEntryListener(this, true);
|
|
||||||
}
|
|
||||||
scheduledJobs = new ConcurrentHashMap<>();
|
|
||||||
updateScheduler();
|
|
||||||
}
|
|
||||||
|
|
||||||
// --- Слушатели Hazelcast Map ---
|
|
||||||
@Override
|
|
||||||
public void entryAdded(EntryEvent<Long, PlannerAllToday> event) {
|
|
||||||
PlannerAllToday task = event.getValue();
|
|
||||||
|
|
||||||
Task taskType = getEnumByKey(Task.class, task.getTask());
|
|
||||||
if (taskType == null) {
|
|
||||||
log.warn("Task skipped, {} task type not recognized", task.getTask());
|
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
processTask(task);
|
if (Active.equalsByKey(taskStatus)) { //&& ACTIVE.equalsById(task.getTaskStatusId())
|
||||||
|
addOrUpdateActiveTask(task);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
public void entryUpdated(PlannerAllToday plannerAllToday) {
|
||||||
public void entryUpdated(EntryEvent<Long, PlannerAllToday> event) {
|
PlannerAllToday task = plannerAllTodayMap.getSingleObjectByID(plannerAllToday.getId());//event.getValue();
|
||||||
PlannerAllToday task = event.getValue();
|
PlannerAllToday oldTask = plannerAllToday;//event.getOldValue();
|
||||||
PlannerAllToday oldTask = event.getOldValue();
|
|
||||||
|
|
||||||
LocalTime oldTime = oldTask.getTaskTime();
|
|
||||||
|
|
||||||
//если таск относится к другому обработчику, пропускаем
|
//если таск относится к другому обработчику, пропускаем
|
||||||
|
|
||||||
if (!Objects.equals(task.getTask(), oldTask.getTask())) {
|
if (!Objects.equals(task.getTask(), oldTask.getTask())) {
|
||||||
throw new IllegalStateException("changed taskId for SchedulerAllToday in core");
|
throw new IllegalStateException("changed taskId for PlannerAllToday in core");
|
||||||
}
|
}
|
||||||
if (Active.equalsByKey(oldTask.getTaskStatus())) { //&& ACTIVE.equalsById(task.getTaskStatusId())
|
String taskStatus = task.getTaskStatus();
|
||||||
if (!removeTask(task.getTask(), oldTime)) {
|
if (taskStatus == null) {
|
||||||
log.debug("cannot cancel task with type {}, time {}", task.getTask(), oldTime.toString());
|
log.warn("Task skipped, {} task status not recognized", task.getTask());
|
||||||
} else {
|
return;
|
||||||
log.debug("task with type {}, time {} execution cancelled, adding altered task...", task.getTask(), oldTime.toString());
|
}
|
||||||
}
|
if (Active.equalsByKey(taskStatus)) { //&& ACTIVE.equalsById(task.getTaskStatusId())
|
||||||
processTask(task);
|
addOrUpdateActiveTask(task);
|
||||||
} else if (Cancel.equalsByKey(oldTask.getTaskStatus())) { //&& TaskStatuses.CANCEL.equalsById(task.getTaskStatusId())
|
} else if (Cancel.equalsByKey(taskStatus)) { //&& TaskStatuses.CANCEL.equalsById(task.getTaskStatusId())
|
||||||
restorePreviouslyRemovedTask(oldTime, oldTask);
|
if (removeTask(task.getTask(), task.getId()))
|
||||||
processTask(task);
|
log.debug("task was successfully deleted after update to cancel");
|
||||||
} else if (Blocked.equalsByKey(oldTask.getTaskStatus())) {
|
else log.debug("no task for deleted after update to cancel");
|
||||||
processTask(task);
|
} else if (Blocked.equalsByKey(taskStatus)) {
|
||||||
|
if (removeTask(task.getTask(), task.getId()))
|
||||||
|
log.debug("task was successfully deleted after update to blocked");
|
||||||
|
else log.debug("no task for deleted after update to blocked");
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
public void entryRemoved(PlannerAllToday taskToRemove) {
|
||||||
public void entryRemoved(EntryEvent<Long, PlannerAllToday> event) {
|
|
||||||
PlannerAllToday taskToRemove = event.getOldValue();
|
|
||||||
if (taskToRemove == null) {
|
if (taskToRemove == null) {
|
||||||
log.debug("no task in removed event");
|
log.debug("no task in removed event");
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
LocalTime removedTaskTime = taskToRemove.getTaskTime();
|
|
||||||
|
|
||||||
if (Active.equalsByKey(taskToRemove.getTaskStatus())) {
|
if (removeTask(taskToRemove.getTask(), taskToRemove.getId()))
|
||||||
if (removeTask(taskToRemove.getTask(), removedTaskTime))
|
log.debug("task successfully removed");
|
||||||
log.debug("task successfully canceled");
|
else log.debug("no task for removing");
|
||||||
} else if (Cancel.equalsByKey(taskToRemove.getTaskStatus())) {
|
|
||||||
restorePreviouslyRemovedTask(removedTaskTime, taskToRemove);
|
|
||||||
} else if (Blocked.equalsByKey(taskToRemove.getTaskStatus())) {
|
|
||||||
log.debug("BLOCKED task removed; do nothing");
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
private void restorePreviouslyRemovedTask(LocalTime oldTime, PlannerAllToday oldTask) {
|
private void restorePreviouslyRemovedTask(Long oldTaskId, PlannerAllToday oldTask) {
|
||||||
ScheduledFuture cancelledFuture = scheduledJobs.get(oldTime);
|
ScheduledFuture cancelledFuture = scheduledJobs.get(oldTaskId);
|
||||||
if (cancelledFuture != null && cancelledFuture.isCancelled()) {
|
if (cancelledFuture != null && cancelledFuture.isCancelled()) {
|
||||||
PlannerAllToday schedulerAllToday = new PlannerAllToday();
|
PlannerAllToday schedulerAllToday = new PlannerAllToday();
|
||||||
schedulerAllToday.setTask(oldTask.getTask());
|
schedulerAllToday.setTask(oldTask.getTask());
|
||||||
schedulerAllToday.setTaskStatus(Active.name());
|
schedulerAllToday.setTaskStatus(Active.name());
|
||||||
schedulerAllToday.setTaskTime(oldTask.getTaskTime());
|
schedulerAllToday.setTaskTime(oldTask.getTaskTime());
|
||||||
log.debug("CANCEL task updated/removed; restoring previously cancelled task");
|
log.debug("CANCEL task updated/removed; restoring previously cancelled task");
|
||||||
processTask(schedulerAllToday);
|
addOrUpdateActiveTask(schedulerAllToday);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -161,49 +172,40 @@ public class TaskManager implements EntryAddedListener<Long, PlannerAllToday>,
|
||||||
schedulerAllTodays.stream().sorted((o1, o2) -> (Active.equalsByKey(o1.getTaskStatus()) && Cancel.equalsByKey(o2.getTaskStatus())) ? -1 : 0)
|
schedulerAllTodays.stream().sorted((o1, o2) -> (Active.equalsByKey(o1.getTaskStatus()) && Cancel.equalsByKey(o2.getTaskStatus())) ? -1 : 0)
|
||||||
.collect(Collectors.toCollection(ArrayList::new));
|
.collect(Collectors.toCollection(ArrayList::new));
|
||||||
for (PlannerAllToday schedulerAllToday : sortedSchedulers) {
|
for (PlannerAllToday schedulerAllToday : sortedSchedulers) {
|
||||||
processTask(schedulerAllToday);
|
addOrUpdateActiveTask(schedulerAllToday);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// --- Работа с задачами ---
|
// --- Работа с задачами ---
|
||||||
private void processTask(PlannerAllToday task) {
|
private void addOrUpdateActiveTask(PlannerAllToday task) {
|
||||||
LocalTime taskTime = task.getTaskTime();
|
LocalTime taskTime = task.getTaskTime();
|
||||||
Task taskType = getEnumByKey(Task.class, task.getTask());
|
Task taskType = getEnumByKey(Task.class, task.getTask());
|
||||||
|
if (taskType == null) {
|
||||||
|
log.warn("Task skipped, {} task type not recognized", task.getTask());
|
||||||
|
return;
|
||||||
|
}
|
||||||
Status taskStatus = getEnumByKey(Status.class, task.getTaskStatus());
|
Status taskStatus = getEnumByKey(Status.class, task.getTaskStatus());
|
||||||
if (taskStatus == null) throw new IllegalStateException("task status from core can't be null");
|
if (taskStatus == null) {
|
||||||
if (taskTime.isBefore(LocalTime.now())) {
|
log.warn("Task skipped, {} task status not recognized", task.getTask());
|
||||||
log.debug("Task skipped - taskTime {} is before now, taskType {} ok", taskTime, task.getTask());
|
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
if (taskStatus.equals(Blocked)) {
|
if (taskStatus.equals(Blocked)) {
|
||||||
log.debug("Task skipped - timeTime {} with status {}", taskTime, Blocked);
|
log.debug("Task skipped - timeTime {} with status {}", taskTime, Blocked);
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
{
|
|
||||||
ScheduledFuture future = scheduledJobs.get(taskTime);
|
|
||||||
if (future != null) {
|
|
||||||
if (Cancel.equals(taskStatus)) {
|
|
||||||
//случай когда пришел cancel, пытаемся отменить зарегистрированный ранее таск
|
|
||||||
future.cancel(false);
|
|
||||||
log.debug("cancelling task time {}", taskTime);
|
|
||||||
return;
|
|
||||||
} else if (Active.equals(taskStatus)) {
|
|
||||||
//случай когда пришел активный таск, и уже был на это время неотмененный
|
|
||||||
if (!future.isCancelled()) {
|
|
||||||
log.debug("such task time {} has already been registered", taskTime);
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
if (Cancel.equals(taskStatus)) {
|
if (Cancel.equals(taskStatus)) {
|
||||||
log.debug("task type {} time {} with status CANCEL - no tasks to cancel found", Objects.requireNonNull(taskType).name(), taskTime);
|
log.debug("task type {} time {} with status CANCEL - no tasks to cancel found", Objects.requireNonNull(taskType).name(), taskTime);
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
//чтобы не получилось, что старый таск не успели отменить а новый уже создали.
|
||||||
|
if (taskTime.isBefore(LocalTime.now().plusSeconds(1))) {
|
||||||
|
log.debug("Task skipped - taskTime {} is before now, taskType {} ok", taskTime, task.getTask());
|
||||||
|
return;
|
||||||
|
}
|
||||||
//пришел активный таск
|
//пришел активный таск
|
||||||
log.debug("adding task type {}, time {}", Objects.requireNonNull(taskType).name(), taskTime);
|
log.debug("adding task type {}, time {}", Objects.requireNonNull(taskType).name(), taskTime);
|
||||||
ScheduledFuture future = taskScheduler.schedule(() -> doJob(task), LocalDateTime.of(LocalDate.now(), taskTime).atZone(ZoneId.systemDefault()).toInstant());
|
ScheduledFuture future = taskScheduler.schedule(() -> doJob(task), LocalDateTime.of(LocalDate.now(), taskTime).atZone(ZoneId.systemDefault()).toInstant());
|
||||||
ScheduledFuture oldFuture = scheduledJobs.put(taskTime, future);
|
ScheduledFuture oldFuture = scheduledJobs.put(task.getId(), future);
|
||||||
if (oldFuture != null && !oldFuture.isCancelled()) { //for synchronization, never
|
if (oldFuture != null && !oldFuture.isCancelled()) { //for synchronization, never
|
||||||
log.warn("tasks were added simultaneously, cancel former one");
|
log.warn("tasks were added simultaneously, cancel former one");
|
||||||
boolean success = oldFuture.cancel(false);
|
boolean success = oldFuture.cancel(false);
|
||||||
|
|
@ -211,10 +213,10 @@ public class TaskManager implements EntryAddedListener<Long, PlannerAllToday>,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private boolean removeTask(String task, LocalTime taskTime) {
|
private boolean removeTask(String task, Long taskId) {
|
||||||
ScheduledFuture future = scheduledJobs.remove(taskTime);
|
ScheduledFuture future = scheduledJobs.remove(taskId);
|
||||||
if (future == null) {
|
if (future == null) {
|
||||||
log.debug("can't cancel task type {}, time {}, not found", task, taskTime);
|
log.debug("can't cancel task type {}, time {}, not found", task, taskId);
|
||||||
return false;
|
return false;
|
||||||
}
|
}
|
||||||
return future.cancel(false);
|
return future.cancel(false);
|
||||||
|
|
@ -230,13 +232,19 @@ public class TaskManager implements EntryAddedListener<Long, PlannerAllToday>,
|
||||||
Instant created = Instant.now();
|
Instant created = Instant.now();
|
||||||
Launcher launcher = new Launcher();
|
Launcher launcher = new Launcher();
|
||||||
launcher.setTask(task.getTask());
|
launcher.setTask(task.getTask());
|
||||||
launcher.setSenderId(task.getParentId());
|
launcher.setSenderId(systemId);
|
||||||
launcher.setCreated(created);
|
launcher.setCreated(created);
|
||||||
launcher.setUpdated(created);
|
launcher.setUpdated(created);
|
||||||
launcherMap.insert(launcher);
|
launcherMap.insert(launcher);
|
||||||
// 2. отправить сообщение
|
// 2. отправить сообщение
|
||||||
launcherSender.sendCommandToQueue(taskE, task.getParentId());
|
launcherSender.sendCommandToQueue(taskE, systemId);
|
||||||
log.debug("successfully processed, new id {}", launcher.getId());
|
log.debug("successfully processed, new id {}", launcher.getId());
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public void close() {
|
||||||
|
log.debug("Closing task manager {}", getClass().getSimpleName());
|
||||||
|
closed.set(true);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -0,0 +1,107 @@
|
||||||
|
package ru.spcex.clearing.scheduler.service;
|
||||||
|
|
||||||
|
import org.junit.jupiter.api.Test;
|
||||||
|
import org.springframework.beans.factory.annotation.Autowired;
|
||||||
|
import ru.clearing.classes.statics.data.scheduler.PlannerAllToday;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
|
||||||
|
import ru.spcex.clearing.scheduler.AbstractServiceTest;
|
||||||
|
import ru.spcex.platform.enumeration.Status;
|
||||||
|
import ru.spcex.platform.enumeration.Task;
|
||||||
|
|
||||||
|
import javax.annotation.PostConstruct;
|
||||||
|
import java.time.LocalTime;
|
||||||
|
|
||||||
|
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.scheduler.config.PlannerQueueConfig.addToPlannerQueue;
|
||||||
|
import static ru.spcex.clearing.scheduler.service.TaskManager.systemId;
|
||||||
|
import static ru.spcex.clearing.test.TestUtils.BASE_REQUEST_MATCHER;
|
||||||
|
import static ru.spcex.clearing.test.config.ImdgTestConfig.currentID;
|
||||||
|
import static ru.spcex.platform.utils.enumeration.IEnumKey.getEnumByKey;
|
||||||
|
|
||||||
|
class TaskManagerTest extends AbstractServiceTest {
|
||||||
|
|
||||||
|
@Autowired
|
||||||
|
LauncherSender launcherSender;
|
||||||
|
|
||||||
|
@PostConstruct
|
||||||
|
public void init() {
|
||||||
|
super.init();
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* {@link TaskManager#entryAdded(PlannerAllToday)}}<br>
|
||||||
|
* Тест проверяет создание и отправку LauncherCommandRequest в Apache Kafka.<br>
|
||||||
|
* Входной запрос {@link PlannerAllToday}:<br>
|
||||||
|
*/
|
||||||
|
@Test
|
||||||
|
void entryAdded() {
|
||||||
|
//ARRANGE
|
||||||
|
PlannerAllToday plannerAllToday = getPlannerAllToday(Task.createOrder.getKey());
|
||||||
|
//ACT
|
||||||
|
addToPlannerQueue(TaskManager.Process.add, plannerAllToday);
|
||||||
|
//ASSERT
|
||||||
|
waitingWhenAddedLauncherCommandRequestAndCheckIt(getEnumByKey(Task.class, plannerAllToday.getTask()), systemId);
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* {@link TaskManager#entryUpdated(PlannerAllToday)}}<br>
|
||||||
|
* Тест проверяет создание и отправку LauncherCommandRequest в Apache Kafka.<br>
|
||||||
|
* Входной запрос {@link PlannerAllToday}:<br>
|
||||||
|
*/
|
||||||
|
@Test
|
||||||
|
void entryUpdated() {
|
||||||
|
//ARRANGE
|
||||||
|
PlannerAllToday plannerAllToday = getPlannerAllToday(Task.createReport_GREP.getKey());
|
||||||
|
PlannerAllToday oldPlanner = getPlannerAllToday(Task.createReport_GREP.getKey());
|
||||||
|
plannerAllToday.setId(oldPlanner.getId());
|
||||||
|
plannerAllTodayImdg.insert(plannerAllToday);
|
||||||
|
addToPlannerQueue(TaskManager.Process.add, oldPlanner);
|
||||||
|
//ACT
|
||||||
|
addToPlannerQueue(TaskManager.Process.update, oldPlanner);
|
||||||
|
//ASSERT
|
||||||
|
waitingWhenAddedLauncherCommandRequestAndCheckIt(getEnumByKey(Task.class, plannerAllToday.getTask()), systemId);
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
/**
|
||||||
|
* {@link TaskManager#entryRemoved(PlannerAllToday)}}<br>
|
||||||
|
* Тест проверяет создание и отправку LauncherCommandRequest в Apache Kafka.<br>
|
||||||
|
* Входной запрос {@link PlannerAllToday}:<br>
|
||||||
|
*/
|
||||||
|
@Test
|
||||||
|
void entryRemoved() {
|
||||||
|
//ARRANGE
|
||||||
|
PlannerAllToday plannerAllToday = getPlannerAllToday(Task.createOrder.getKey());
|
||||||
|
PlannerAllToday oldPlanner = getPlannerAllToday(Task.createReport_GREP.getKey());
|
||||||
|
addToPlannerQueue(TaskManager.Process.add, oldPlanner);
|
||||||
|
addToPlannerQueue(TaskManager.Process.delete, oldPlanner);
|
||||||
|
//ACT
|
||||||
|
addToPlannerQueue(TaskManager.Process.add, plannerAllToday);
|
||||||
|
//ASSERT
|
||||||
|
waitingWhenAddedLauncherCommandRequestAndCheckIt(getEnumByKey(Task.class, plannerAllToday.getTask()), systemId);
|
||||||
|
}
|
||||||
|
|
||||||
|
protected PlannerAllToday getPlannerAllToday(String task) {
|
||||||
|
PlannerAllToday plannerAllToday = new PlannerAllToday();
|
||||||
|
plannerAllToday.setId(currentID.getAndIncrement());
|
||||||
|
plannerAllToday.setTask(task);
|
||||||
|
plannerAllToday.setTaskTime(LocalTime.now().plusSeconds(18));
|
||||||
|
plannerAllToday.setTaskStatus(Status.Active.getKey());
|
||||||
|
return plannerAllToday;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void waitingWhenAddedLauncherCommandRequestAndCheckIt(Task toTaskQueue, Long userId) {
|
||||||
|
BaseRequest<Object> predictableBaseRequest = launcherSender.makeCmdRequest(toTaskQueue, userId);
|
||||||
|
|
||||||
|
//waiting for kafka producer send message (finale event)
|
||||||
|
verify(mockProducer, timeout(60_000L).times(1))
|
||||||
|
.send(producerRecord.capture());
|
||||||
|
|
||||||
|
BaseRequest<Object> baseRequestResult = (BaseRequest<Object>) producerRecord.getValue().value();
|
||||||
|
assertEquals(toTaskQueue.topic(), producerRecord.getValue().topic());
|
||||||
|
predictableBaseRequest.setId(baseRequestResult.getId());
|
||||||
|
BASE_REQUEST_MATCHER.assertMatch(baseRequestResult, predictableBaseRequest);
|
||||||
|
}
|
||||||
|
}
|
||||||
Loading…
Add table
Reference in a new issue