From f41aecf7391a0665ba54fbe0351afa43440baa8f Mon Sep 17 00:00:00 2001 From: psemenkov Date: Mon, 24 Apr 2023 17:57:34 +0300 Subject: [PATCH] =?UTF-8?q?http://jira.mfd.msk:8088/browse/CLS-268=20?= =?UTF-8?q?=D0=A1=D0=B4=D0=B5=D0=BB=D0=B0=D0=BB=20TaskManager=20=D0=B1?= =?UTF-8?q?=D0=B5=D0=B7=20hazelcast=20listener=20=D0=B4=D0=BE=D0=B1=D0=B0?= =?UTF-8?q?=D0=B2=D0=B8=D0=BB=20=D0=BE=D1=82=D0=BF=D1=80=D0=B0=D0=B2=D0=BA?= =?UTF-8?q?=D1=83=20=D1=81=D0=BE=D0=BE=D0=B1=D1=89=D0=B5=D0=BD=D0=B8=D0=B9?= =?UTF-8?q?=20=D0=B8=D0=B7=20=D1=81=D0=B5=D1=80=D0=B2=D0=B8=D1=81=D0=BE?= =?UTF-8?q?=D0=B2=20=D0=B2=20=D0=BE=D1=87=D0=B5=D1=80=D0=B5=D0=B4=D1=8C=20?= =?UTF-8?q?=D0=B8=20=D0=BD=D0=B0=D0=BF=D0=B8=D1=81=D0=B0=D0=BB=20TaskManag?= =?UTF-8?q?erTest.?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../scheduler/config/PlannerQueueConfig.java | 29 +++ .../service/ClearingCalendarService.java | 28 +-- .../scheduler/service/LauncherSender.java | 1 + .../scheduler/service/PlannerService.java | 20 +- .../service/PlannerTemplateService.java | 26 ++- .../scheduler/service/TaskManager.java | 208 +++++++++--------- .../scheduler/service/TaskManagerTest.java | 107 +++++++++ 7 files changed, 293 insertions(+), 126 deletions(-) create mode 100644 clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/config/PlannerQueueConfig.java create mode 100644 clearing-parent/scheduler-service/src/test/java/ru/spcex/clearing/scheduler/service/TaskManagerTest.java diff --git a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/config/PlannerQueueConfig.java b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/config/PlannerQueueConfig.java new file mode 100644 index 000000000..0534b2053 --- /dev/null +++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/config/PlannerQueueConfig.java @@ -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> plannerQueue; + + @Bean(name = "plannerQueue") + public BlockingQueue> 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) { + } + } +} 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 04a46c4b8..d53bec1bb 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 @@ -9,7 +9,6 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.stereotype.Service; 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.PlannerTemplate; 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.scheduler.PlannerAllTodayBuilder; 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.Status; +import ru.spcex.platform.enumeration.Parent; import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.ImdgProvider; import ru.spcex.platform.utils.enumeration.IMessageResolver; @@ -35,11 +33,11 @@ import java.time.Instant; import java.time.LocalDate; import java.util.Arrays; import java.util.Collection; -import java.util.List; import java.util.Map; import java.util.function.Function; import static ru.spcex.clearing.scheduler.ISchedulerChecker.isValidWorkday; +import static ru.spcex.clearing.scheduler.config.PlannerQueueConfig.addToPlannerQueue; @Service public class ClearingCalendarService extends QueueConsumer implements InitializingBean { @@ -47,7 +45,6 @@ public class ClearingCalendarService extends QueueConsumer implements Initializi private final Imdg clearingCalendarMap; private final Imdg plannerTemplateMap; private final Imdg plannerAllTodayMap; - private final Imdg plannerMap; private final IMessageResolver messageResolver; private final UserRoleVerification userRoleVerification; private final Function clearingCalendarDeleteRequestValidation; @@ -65,7 +62,6 @@ public class ClearingCalendarService extends QueueConsumer implements Initializi super(kafkaQueue, kafkaProducer); this.clearingCalendarMap = imdgProvider.getImdg(IMDGDistributedNames.Map_ClearingCalendar, ClearingCalendar.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.userRoleVerification = userRoleVerification; this.clearingCalendarDeleteRequestValidation = clearingCalendarDeleteRequestValidator; @@ -157,8 +153,11 @@ public class ClearingCalendarService extends QueueConsumer implements Initializi if (isValidWorkday(clearingCalendar, isWeekend)) { for (PlannerTemplate plannerTemplate : plannerTemplateMap.getAllValues()) { PlannerAllToday plannerAllToday = plannerAllTodayMap.getSingleObjectBySQL(String.format("parentId = %s", plannerTemplate.getId())); - if (plannerAllToday == null) - plannerAllTodayMap.insert(PlannerAllTodayBuilder.builder().append(plannerTemplate).build()); + if (plannerAllToday == null) { + 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 - List plannerIds = plannerMap.getCollectionObjectsByFieldValues(Map.of("taskStatus", Status.Active.getKey(), - "clearingDate", currentDate)).stream() - .mapToLong(SpcexObjectBase::getId) - .boxed().toList(); + Collection values = plannerAllTodayMap.getCollectionObjectsByFieldValues(Map.of("parent", Parent.Template.getKey())); - Collection values = plannerAllTodayMap.getAllValues().stream() - .filter(plannerAllToday -> !plannerIds.contains(plannerAllToday.getParentId())).toList(); - - values.forEach(plannerAllTodayMap::delete); + values.forEach(plannerAllToday -> { + plannerAllTodayMap.delete(plannerAllToday); + addToPlannerQueue(TaskManager.Process.delete, plannerAllToday); + }); } } diff --git a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/LauncherSender.java b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/LauncherSender.java index 0f7d385ee..e58529cc9 100644 --- a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/LauncherSender.java +++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/LauncherSender.java @@ -39,6 +39,7 @@ public class LauncherSender { BaseRequest request = new BaseRequest<>(); request.setId(idGenerator.nextId()); request.setActionType(ActionType.NEW); + request.setUserId(userId); request.setRequestPayload(toRequest);// iAction.toRequest() return request; } diff --git a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/PlannerService.java b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/PlannerService.java index dddcc9ba7..6376777c2 100644 --- a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/PlannerService.java +++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/PlannerService.java @@ -32,6 +32,8 @@ import java.util.Collection; import java.util.Map; import java.util.function.Function; +import static ru.spcex.clearing.scheduler.config.PlannerQueueConfig.addToPlannerQueue; + @Service public class PlannerService extends QueueConsumer implements InitializingBean { private final Logger log = LoggerFactory.getLogger(getClass()); @@ -151,9 +153,16 @@ public class PlannerService extends QueueConsumer implements InitializingBean { public void cudPlannerAllToday(Planner planner) { if (planner.getTaskStatus().equalsIgnoreCase(Status.Active.getKey())) { PlannerAllToday plannerAllToday = PlannerAllTodayBuilder.builder().append(planner).build(); - PlannerAllToday allToday = plannerAllTodayMap.getSingleObjectBySQL(String.format("parentId = %s", planner.getId())); - if (allToday != null) plannerAllToday.setId(allToday.getId()); - plannerAllTodayMap.insert(plannerAllToday); + PlannerAllToday oldPlannerAllToday = plannerAllTodayMap.getSingleObjectBySQL(String.format("parentId = %s", planner.getId())); + if (oldPlannerAllToday != null) { + 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())) { @@ -162,7 +171,10 @@ public class PlannerService extends QueueConsumer implements InitializingBean { "taskTime", planner.getTaskTime(), "securityId", planner.getSecurityId(), "companyId", planner.getCompanyId())); - values.forEach(plannerAllTodayMap::delete); + values.forEach(plannerAllToday -> { + plannerAllTodayMap.delete(plannerAllToday); + addToPlannerQueue(TaskManager.Process.delete, plannerAllToday); + }); } } } 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 443f0c540..b5d1fa993 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 @@ -36,6 +36,7 @@ import java.util.Map; import java.util.function.Function; import static ru.spcex.clearing.scheduler.ISchedulerChecker.isValidWorkday; +import static ru.spcex.clearing.scheduler.config.PlannerQueueConfig.addToPlannerQueue; @Service public class PlannerTemplateService extends QueueConsumer implements InitializingBean { @@ -102,7 +103,11 @@ public class PlannerTemplateService extends QueueConsumer implements Initializin plannerTemplate.setSecurityId(req.getSecurityId()); 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()); return null; } @@ -127,10 +132,16 @@ public class PlannerTemplateService extends QueueConsumer implements Initializin plannerTemplateMap.update(plannerTemplate); if (todayWorkDay()) { - PlannerAllToday planner = PlannerAllTodayBuilder.builder().append(plannerTemplate).build(); - PlannerAllToday plannerAllToday = plannerAllTodayMap.getSingleObjectBySQL(String.format("parentId = %s", plannerTemplate.getId())); - if (plannerAllToday != null) planner.setId(plannerAllToday.getId()); - plannerAllTodayMap.insert(planner); + PlannerAllToday plannerAllToday = PlannerAllTodayBuilder.builder().append(plannerTemplate).build(); + PlannerAllToday oldPlannerAllToday = plannerAllTodayMap.getSingleObjectBySQL(String.format("parentId = %s", plannerTemplate.getId())); + if (oldPlannerAllToday != null) { + 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()); return null; @@ -149,7 +160,10 @@ public class PlannerTemplateService extends QueueConsumer implements Initializin if (todayWorkDay()) { 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); return null; diff --git a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/TaskManager.java b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/TaskManager.java index bcc7fbc07..f4c543bc2 100644 --- a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/TaskManager.java +++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/TaskManager.java @@ -1,14 +1,10 @@ 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.LoggerFactory; import org.springframework.beans.factory.InitializingBean; 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.scheduling.TaskScheduler; 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.imdg.api.Imdg; import ru.spcex.platform.imdg.api.ImdgProvider; -import ru.spcex.platform.imdg.iml.hazelcast.adapter.ImdgHazelcast; import java.time.*; -import java.util.ArrayList; -import java.util.Collection; -import java.util.Date; -import java.util.Objects; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.ScheduledFuture; +import java.util.*; +import java.util.concurrent.*; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.function.Consumer; import java.util.stream.Collectors; import static ru.spcex.clearing.imdg.IMDGDistributedNames.Map_Launcher; @@ -39,25 +32,63 @@ import static ru.spcex.platform.utils.enumeration.IEnumKey.getEnumByKey; *

*/ @Service("userReportTask") -@Lazy -public class TaskManager implements EntryAddedListener, - EntryUpdatedListener, EntryRemovedListener, - InitializingBean { +public class TaskManager implements InitializingBean, AutoCloseable { private static final Logger log = LoggerFactory.getLogger(TaskManager.class); + public static final Long systemId = 0L;//для системных задач по договоренности, должен быть пользователь с правами админа и id=0L + private final TaskScheduler taskScheduler; private final ImdgProvider imdgProvider; private Imdg launcherMap; private Imdg plannerAllTodayMap; - private ConcurrentHashMap scheduledJobs; + private ConcurrentHashMap scheduledJobs; private LauncherSender launcherSender; + private final BlockingQueue> plannerQueue; + private final Map> callbacks = new HashMap<>(); + private final AtomicBoolean closed = new AtomicBoolean(false); + private final ExecutorService outputExecutor; @Autowired TaskManager(TaskScheduler taskScheduler, ImdgProvider imdgProvider, - LauncherSender launcherSender) { + LauncherSender launcherSender, + @Qualifier("plannerQueue") BlockingQueue> plannerQueue) { this.taskScheduler = taskScheduler; this.imdgProvider = imdgProvider; 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 entry = plannerQueue.take(); + callbacks.get(entry.getKey()).accept(entry.getValue()); + } + } catch (InterruptedException e) { +// closed.set(true); + } + }); } private static LocalDateTime dateOldTypeConvert(@NonNull Date oldDate) { @@ -72,85 +103,65 @@ public class TaskManager implements EntryAddedListener, return dateOldTypeConvert(oldDate).toLocalTime(); } - @Override - public void afterPropertiesSet() { - this.launcherMap = imdgProvider.getImdg(Map_Launcher, Launcher.class); - this.plannerAllTodayMap = imdgProvider.getImdg(Map_PlannerAllToday, PlannerAllToday.class); - if (plannerAllTodayMap instanceof ImdgHazelcast plannerAllTodayImdgHazelcast) { - plannerAllTodayImdgHazelcast.getMap().addEntryListener(this, true); - } - scheduledJobs = new ConcurrentHashMap<>(); - updateScheduler(); - } - - // --- Слушатели Hazelcast Map --- - @Override - public void entryAdded(EntryEvent 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()); + // --- Обработчики PlannerAllToday --- + public void entryAdded(PlannerAllToday task) { + String taskStatus = task.getTaskStatus(); + if (taskStatus == null) { + log.warn("Task skipped, {} task status not recognized", task.getTask()); return; } - processTask(task); + if (Active.equalsByKey(taskStatus)) { //&& ACTIVE.equalsById(task.getTaskStatusId()) + addOrUpdateActiveTask(task); + } } - @Override - public void entryUpdated(EntryEvent event) { - PlannerAllToday task = event.getValue(); - PlannerAllToday oldTask = event.getOldValue(); - - LocalTime oldTime = oldTask.getTaskTime(); + public void entryUpdated(PlannerAllToday plannerAllToday) { + PlannerAllToday task = plannerAllTodayMap.getSingleObjectByID(plannerAllToday.getId());//event.getValue(); + PlannerAllToday oldTask = plannerAllToday;//event.getOldValue(); //если таск относится к другому обработчику, пропускаем 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()) - if (!removeTask(task.getTask(), oldTime)) { - log.debug("cannot cancel task with type {}, time {}", task.getTask(), oldTime.toString()); - } else { - log.debug("task with type {}, time {} execution cancelled, adding altered task...", task.getTask(), oldTime.toString()); - } - processTask(task); - } else if (Cancel.equalsByKey(oldTask.getTaskStatus())) { //&& TaskStatuses.CANCEL.equalsById(task.getTaskStatusId()) - restorePreviouslyRemovedTask(oldTime, oldTask); - processTask(task); - } else if (Blocked.equalsByKey(oldTask.getTaskStatus())) { - processTask(task); + String taskStatus = task.getTaskStatus(); + if (taskStatus == null) { + log.warn("Task skipped, {} task status not recognized", task.getTask()); + return; + } + if (Active.equalsByKey(taskStatus)) { //&& ACTIVE.equalsById(task.getTaskStatusId()) + addOrUpdateActiveTask(task); + } else if (Cancel.equalsByKey(taskStatus)) { //&& TaskStatuses.CANCEL.equalsById(task.getTaskStatusId()) + if (removeTask(task.getTask(), task.getId())) + log.debug("task was successfully deleted after update to cancel"); + else log.debug("no task for deleted after update to cancel"); + } 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(EntryEvent event) { - PlannerAllToday taskToRemove = event.getOldValue(); + public void entryRemoved(PlannerAllToday taskToRemove) { if (taskToRemove == null) { log.debug("no task in removed event"); return; } - LocalTime removedTaskTime = taskToRemove.getTaskTime(); - if (Active.equalsByKey(taskToRemove.getTaskStatus())) { - if (removeTask(taskToRemove.getTask(), removedTaskTime)) - log.debug("task successfully canceled"); - } else if (Cancel.equalsByKey(taskToRemove.getTaskStatus())) { - restorePreviouslyRemovedTask(removedTaskTime, taskToRemove); - } else if (Blocked.equalsByKey(taskToRemove.getTaskStatus())) { - log.debug("BLOCKED task removed; do nothing"); - } + if (removeTask(taskToRemove.getTask(), taskToRemove.getId())) + log.debug("task successfully removed"); + else log.debug("no task for removing"); } - private void restorePreviouslyRemovedTask(LocalTime oldTime, PlannerAllToday oldTask) { - ScheduledFuture cancelledFuture = scheduledJobs.get(oldTime); + private void restorePreviouslyRemovedTask(Long oldTaskId, PlannerAllToday oldTask) { + ScheduledFuture cancelledFuture = scheduledJobs.get(oldTaskId); if (cancelledFuture != null && cancelledFuture.isCancelled()) { PlannerAllToday schedulerAllToday = new PlannerAllToday(); schedulerAllToday.setTask(oldTask.getTask()); schedulerAllToday.setTaskStatus(Active.name()); schedulerAllToday.setTaskTime(oldTask.getTaskTime()); log.debug("CANCEL task updated/removed; restoring previously cancelled task"); - processTask(schedulerAllToday); + addOrUpdateActiveTask(schedulerAllToday); } } @@ -161,49 +172,40 @@ public class TaskManager implements EntryAddedListener, schedulerAllTodays.stream().sorted((o1, o2) -> (Active.equalsByKey(o1.getTaskStatus()) && Cancel.equalsByKey(o2.getTaskStatus())) ? -1 : 0) .collect(Collectors.toCollection(ArrayList::new)); for (PlannerAllToday schedulerAllToday : sortedSchedulers) { - processTask(schedulerAllToday); + addOrUpdateActiveTask(schedulerAllToday); } } // --- Работа с задачами --- - private void processTask(PlannerAllToday task) { + private void addOrUpdateActiveTask(PlannerAllToday task) { LocalTime taskTime = task.getTaskTime(); 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()); - if (taskStatus == null) throw new IllegalStateException("task status from core can't be null"); - if (taskTime.isBefore(LocalTime.now())) { - log.debug("Task skipped - taskTime {} is before now, taskType {} ok", taskTime, task.getTask()); + if (taskStatus == null) { + log.warn("Task skipped, {} task status not recognized", task.getTask()); return; } if (taskStatus.equals(Blocked)) { log.debug("Task skipped - timeTime {} with status {}", taskTime, Blocked); 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)) { log.debug("task type {} time {} with status CANCEL - no tasks to cancel found", Objects.requireNonNull(taskType).name(), taskTime); 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); 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 log.warn("tasks were added simultaneously, cancel former one"); boolean success = oldFuture.cancel(false); @@ -211,10 +213,10 @@ public class TaskManager implements EntryAddedListener, } } - private boolean removeTask(String task, LocalTime taskTime) { - ScheduledFuture future = scheduledJobs.remove(taskTime); + private boolean removeTask(String task, Long taskId) { + ScheduledFuture future = scheduledJobs.remove(taskId); 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 future.cancel(false); @@ -230,13 +232,19 @@ public class TaskManager implements EntryAddedListener, Instant created = Instant.now(); Launcher launcher = new Launcher(); launcher.setTask(task.getTask()); - launcher.setSenderId(task.getParentId()); + launcher.setSenderId(systemId); launcher.setCreated(created); launcher.setUpdated(created); launcherMap.insert(launcher); // 2. отправить сообщение - launcherSender.sendCommandToQueue(taskE, task.getParentId()); + launcherSender.sendCommandToQueue(taskE, systemId); log.debug("successfully processed, new id {}", launcher.getId()); } } + + @Override + public void close() { + log.debug("Closing task manager {}", getClass().getSimpleName()); + closed.set(true); + } } diff --git a/clearing-parent/scheduler-service/src/test/java/ru/spcex/clearing/scheduler/service/TaskManagerTest.java b/clearing-parent/scheduler-service/src/test/java/ru/spcex/clearing/scheduler/service/TaskManagerTest.java new file mode 100644 index 000000000..18cbebc9a --- /dev/null +++ b/clearing-parent/scheduler-service/src/test/java/ru/spcex/clearing/scheduler/service/TaskManagerTest.java @@ -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)}}
+ * Тест проверяет создание и отправку LauncherCommandRequest в Apache Kafka.
+ * Входной запрос {@link PlannerAllToday}:
+ */ + @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)}}
+ * Тест проверяет создание и отправку LauncherCommandRequest в Apache Kafka.
+ * Входной запрос {@link PlannerAllToday}:
+ */ + @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)}}
+ * Тест проверяет создание и отправку LauncherCommandRequest в Apache Kafka.
+ * Входной запрос {@link PlannerAllToday}:
+ */ + @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 predictableBaseRequest = launcherSender.makeCmdRequest(toTaskQueue, userId); + + //waiting for kafka producer send message (finale event) + verify(mockProducer, timeout(60_000L).times(1)) + .send(producerRecord.capture()); + + BaseRequest baseRequestResult = (BaseRequest) producerRecord.getValue().value(); + assertEquals(toTaskQueue.topic(), producerRecord.getValue().topic()); + predictableBaseRequest.setId(baseRequestResult.getId()); + BASE_REQUEST_MATCHER.assertMatch(baseRequestResult, predictableBaseRequest); + } +} \ No newline at end of file