clearingCalendars = clearingCalendarMap.getCollectionObjectsByFieldValues(Map.of("dayStatus", DayStatus.Workday.getKey(),
+ "clearingDate", currentDate));
+ if (clearingCalendars != null && !clearingCalendars.isEmpty()) {
+ for (ClearingCalendar calendar : clearingCalendars) {
+ if (isValidWorkday(calendar, isWeekend)) return true;
+ }
+ }
+ return false;
+ }
}
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..9aed3335e 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,15 +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.lang.NonNull;
+import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.scheduling.TaskScheduler;
import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.scheduler.Launcher;
@@ -18,15 +13,13 @@ 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 ru.spcex.platform.utils.log.ExceptionUtils;
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,118 +32,127 @@ 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();
}
- private static LocalDateTime dateOldTypeConvert(@NonNull Date oldDate) {
- return LocalDateTime.ofInstant(oldDate.toInstant(), ZoneId.systemDefault());
- }
-
- private static LocalDate dateTypeConvert(@NonNull Date oldDate) {
- return dateOldTypeConvert(oldDate).toLocalDate();
- }
-
- private static LocalTime timeTypeConvert(@NonNull Date oldDate) {
- return dateOldTypeConvert(oldDate).toLocalTime();
+ 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);
- if (plannerAllTodayMap instanceof ImdgHazelcast plannerAllTodayImdgHazelcast) {
- plannerAllTodayImdgHazelcast.getMap().addEntryListener(this, true);
- }
+ 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();
}
- // --- Слушатели Hazelcast Map ---
- @Override
- public void entryAdded(EntryEvent event) {
- PlannerAllToday task = event.getValue();
+ private void process() {
+ outputExecutor.submit(() -> {
+ while (!closed.get()) {
+ try {
+ //процесс ожидает пока появится новое сообщение в очереди
+ Map.Entry entry = plannerQueue.take();
+ callbacks.get(entry.getKey()).accept(entry.getValue());
+ } catch (InterruptedException e) {
+ closed.set(true);
+ Thread.currentThread().interrupt();
+ } catch (Throwable t) {
+ log.error("InterruptedException in process, {}", ExceptionUtils.getStackTrace(t));
+ }
+ }
+ });
+ }
- 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 +163,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 +204,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 +223,20 @@ 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);
+ outputExecutor.shutdown();
+ }
}
diff --git a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/validation/rules/common/TimeNotBeforeRule.java b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/validation/rules/common/TimeNotBeforeRule.java
index 90287864e..acdc2df60 100644
--- a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/validation/rules/common/TimeNotBeforeRule.java
+++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/validation/rules/common/TimeNotBeforeRule.java
@@ -11,27 +11,42 @@ import java.util.function.Function;
/**
* Проверка поля с LocalTime. Условие проверки: проверяемое время >= текущее время
+ *
* @param Класс проверяемого объекта
*/
-public record TimeNotBeforeRule(String fieldName, Function getter, boolean required) implements IValidationRule> {
+public record TimeNotBeforeRule(String fieldName,
+ Function getter,
+ boolean required,
+ boolean isToday) implements IValidationRule> {
/**
* @param fieldName Название поля класса, используется для передачи ошибки
- * @param getter Метод получения проверяемого времени
- * @param required Флаг обязательности поля
- * @param Класс проверяемого объекта
+ * @param getter Метод получения проверяемого времени
+ * @param required Флаг обязательности поля
+ * @param Класс проверяемого объекта
+ * @param isToday Флаг что дата сегодняшняя
*/
- public static TimeNotBeforeRule instance(String fieldName, Function getter, boolean required) {
- return new TimeNotBeforeRule<>(fieldName, getter, required);
+ public static TimeNotBeforeRule instance(String fieldName, Function getter, boolean required, boolean isToday) {
+ return new TimeNotBeforeRule<>(fieldName, getter, required, isToday);
}
/**
* @param fieldName Название поля класса, используется для передачи ошибки
- * @param getter Метод получения проверяемого времени
- * @param Класс проверяемого объекта
+ * @param getter Метод получения проверяемого времени
+ * @param required Флаг обязательности поля
+ * @param Класс проверяемого объекта
+ */
+ public static TimeNotBeforeRule instance(String fieldName, Function getter, boolean required) {
+ return new TimeNotBeforeRule<>(fieldName, getter, required, true);
+ }
+
+ /**
+ * @param fieldName Название поля класса, используется для передачи ошибки
+ * @param getter Метод получения проверяемого времени
+ * @param Класс проверяемого объекта
*/
public static TimeNotBeforeRule instance(String fieldName, Function getter) {
- return new TimeNotBeforeRule<>(fieldName, getter, true);
+ return new TimeNotBeforeRule<>(fieldName, getter, true, true);
}
@Override
@@ -39,7 +54,8 @@ public record TimeNotBeforeRule(String fieldName, Function gett
R validatedObject = context.getValidatedObject();
LocalTime date = getter.apply(validatedObject);
if (date == null) return required ? of(ValidationError.EmptyRequiredValue, fieldName) : Optional.empty();
- if (date.isBefore(LocalTime.now())) return of(ValidationError.TaskForPastTime, fieldName);
+ if (date.isBefore(LocalTime.now()))
+ return isToday ? of(ValidationError.TaskForPastTime, fieldName) : Optional.empty();
return Optional.empty();
}
}
diff --git a/clearing-parent/scheduler-service/src/test/java/ru/spcex/clearing/scheduler/AbstractServiceTest.java b/clearing-parent/scheduler-service/src/test/java/ru/spcex/clearing/scheduler/AbstractServiceTest.java
new file mode 100644
index 000000000..c2a6b2173
--- /dev/null
+++ b/clearing-parent/scheduler-service/src/test/java/ru/spcex/clearing/scheduler/AbstractServiceTest.java
@@ -0,0 +1,183 @@
+package ru.spcex.clearing.scheduler;
+
+import org.apache.kafka.clients.producer.MockProducer;
+import org.apache.kafka.clients.producer.ProducerRecord;
+import org.junit.jupiter.api.extension.ExtendWith;
+import org.mockito.ArgumentCaptor;
+import org.mockito.Captor;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.beans.factory.annotation.Qualifier;
+import org.springframework.boot.test.mock.mockito.MockBean;
+import org.springframework.test.context.ContextConfiguration;
+import org.springframework.test.context.junit.jupiter.SpringExtension;
+import ru.clearing.classes.statics.data.company.Company;
+import ru.clearing.classes.statics.data.scheduler.*;
+import ru.clearing.classes.statics.data.security.Security;
+import ru.clearing.platform.dictionary.DayStatusDictionary;
+import ru.clearing.platform.dictionary.TaskDictionary;
+import ru.clearing.platform.dictionary.TaskStatusDictionary;
+import ru.spcex.clearing.imdg.IMDGDistributedNames;
+import ru.spcex.clearing.scheduler.config.ErrorResolverConfig;
+import ru.spcex.clearing.scheduler.config.PlannerQueueConfig;
+import ru.spcex.clearing.scheduler.config.SchedulerTestConfig;
+import ru.spcex.clearing.scheduler.config.validation.ClearingCalendarValidationConfig;
+import ru.spcex.clearing.scheduler.config.validation.PlannerTemplateValidationConfig;
+import ru.spcex.clearing.scheduler.config.validation.PlannerValidationConfig;
+import ru.spcex.clearing.scheduler.config.validation.ValidationConfig;
+import ru.spcex.clearing.scheduler.service.*;
+import ru.spcex.clearing.test.MatcherFactory;
+import ru.spcex.clearing.test.TestUtils;
+import ru.spcex.clearing.test.config.ImdgTestConfig;
+import ru.spcex.clearing.test.config.KafkaTestConfig;
+import ru.spcex.platform.enumeration.DayStatus;
+import ru.spcex.platform.enumeration.Status;
+import ru.spcex.platform.enumeration.Task;
+import ru.spcex.platform.enumeration.WorkflowStatus;
+import ru.spcex.platform.imdg.api.Imdg;
+import ru.spcex.platform.imdg.api.ImdgProvider;
+
+import java.time.LocalDate;
+import java.time.LocalTime;
+
+import static org.mockito.Mockito.doReturn;
+import static org.mockito.Mockito.spy;
+import static ru.spcex.clearing.test.MatcherFactory.usingIgnoringFieldsComparator;
+import static ru.spcex.clearing.test.config.ImdgTestConfig.currentID;
+import static ru.spcex.clearing.test.config.ImdgTestConfig.waitAvailableImdgProviderAndAddAdminWithDefaultId;
+import static ru.spcex.platform.enumeration.Market.mkrs;
+
+@ExtendWith(SpringExtension.class)
+@ContextConfiguration(classes = {
+ PlannerService.class,
+ PlannerTemplateValidationConfig.class,
+ ClearingCalendarService.class,
+ PlannerTemplateService.class,
+ TaskManager.class,
+ LauncherSender.class,
+ PlannerValidationConfig.class,
+ ClearingCalendarValidationConfig.class,
+ ValidationConfig.class,
+ LauncherService.class,
+ ErrorResolverConfig.class,
+ PlannerQueueConfig.class,
+ SchedulerTestConfig.class,
+ ImdgTestConfig.class,
+ KafkaTestConfig.class})
+public abstract class AbstractServiceTest {
+ protected static final MatcherFactory.Matcher PLANNER_ALL_TODAY_MATCHER = usingIgnoringFieldsComparator("created", "updated");
+ protected static final long id = currentID.getAndIncrement();
+ protected Imdg plannerAllTodayImdg;
+ protected Imdg plannerTemplateImdg;
+ protected Imdg clearingCalendarImdg;
+ protected Imdg plannerImdg;
+ protected Imdg launcherMap;
+ protected long newCompanyId = currentID.getAndIncrement();
+ protected long updateCompanyId = currentID.getAndIncrement();
+ protected long deleteCompanyId = currentID.getAndIncrement();
+ protected long testSecurityId = currentID.getAndIncrement();
+
+ @Captor
+ protected ArgumentCaptor producerRecord;
+ @MockBean
+ protected MockProducer mockProducer;
+ @Autowired
+ @Qualifier("hazelcastServiceTest")
+ protected ImdgProvider imdgProvider;
+
+ protected void init() {
+ waitAvailableImdgProviderAndAddAdminWithDefaultId();
+ this.plannerAllTodayImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_PlannerAllToday, PlannerAllToday.class);
+ this.plannerTemplateImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_PlannerTemplate, PlannerTemplate.class);
+ this.clearingCalendarImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ClearingCalendar, ClearingCalendar.class);
+ this.plannerImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Planner, Planner.class);
+ this.launcherMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Launcher, Launcher.class);
+
+ TaskDictionary taskDictionary = new TaskDictionary();
+ Imdg taskDictionaryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_TaskDictionary, TaskDictionary.class);
+ taskDictionary.setCode(Task.accountBlock.getKey());
+ taskDictionary.setId(id);
+ taskDictionaryImdg.insert(taskDictionary);
+
+ DayStatusDictionary dayStatusDictionary = new DayStatusDictionary();
+ dayStatusDictionary.setId(0L);
+ dayStatusDictionary.setCode(DayStatus.Workday.getKey());
+ Imdg dayStatusDictionaryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_DayStatusDictionary, DayStatusDictionary.class);
+ dayStatusDictionaryImdg.insert(dayStatusDictionary);
+
+ TaskStatusDictionary taskStatusDictionary = new TaskStatusDictionary();
+ taskStatusDictionary.setCode(Status.Active.getKey());
+ taskStatusDictionary.setId(id);
+ Imdg taskStatusDictionaryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_TaskStatusDictionary, TaskStatusDictionary.class);
+ taskStatusDictionaryImdg.insert(taskStatusDictionary);
+
+ Company company = new Company();
+ company.setId(newCompanyId);
+ company.setWorkflowStatus(WorkflowStatus.Active.getKey());
+ Imdg companyImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Company, Company.class);
+ companyImdg.insert(company);
+ company.setId(updateCompanyId);
+ company.setWorkflowStatus(WorkflowStatus.Active.getKey());
+ companyImdg.insert(company);
+ company.setId(deleteCompanyId);
+ company.setWorkflowStatus(WorkflowStatus.Active.getKey());
+ companyImdg.insert(company);
+
+ Security security = new Security();
+ security.setId(testSecurityId);
+ security.setWorkflowStatus(WorkflowStatus.Active.getKey());
+ Imdg securityImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Security, Security.class);
+ securityImdg.insert(security);
+
+ clearingCalendarImdg.insert(getClearingCalendar());
+
+ TestUtils.FutureRecordMetadata future = spy(new TestUtils.FutureRecordMetadata());
+ doReturn(future).when(mockProducer).send(producerRecord.capture());
+ }
+
+ protected PlannerTemplate getPlannerTemplate(long companyId) {
+ PlannerTemplate plannerTemplate = new PlannerTemplate();
+ plannerTemplate.setId(currentID.getAndIncrement());
+ plannerTemplate.setTask(Task.accountBlock.getKey());
+ plannerTemplate.setTaskTime(LocalTime.now().plusHours(1).withNano(0));
+ plannerTemplate.setTaskStatus(Status.Active.getKey());
+ plannerTemplate.setCompanyId(companyId);
+ plannerTemplate.setSecurityId(testSecurityId);
+ return plannerTemplate;
+ }
+
+ protected Planner getPlanner(long companyId) {
+ Planner planner = new Planner();
+ planner.setId(currentID.getAndIncrement());
+ planner.setTask(Task.accountBlock.getKey());
+ planner.setTaskTime(LocalTime.now().plusHours(1).withNano(0));
+ planner.setClearingDate(LocalDate.now());
+ planner.setMarket(mkrs.getKey());
+ planner.setTaskStatus(Status.Active.getKey());
+ planner.setCompanyId(companyId);
+ planner.setSecurityId(testSecurityId);
+ return planner;
+ }
+
+ protected void checkPlannerAllTodayByPlannerTemplate(PlannerTemplate plannerTemplate) {
+ PlannerAllToday plannerAllToday = PlannerAllTodayBuilder.builder().append(plannerTemplate).build();
+ PlannerAllToday plannerAllTodayRes = plannerAllTodayImdg.getSingleObjectBySQL(String.format("companyId = %s", plannerTemplate.getCompanyId()));
+ plannerAllToday.setId(plannerAllTodayRes.getId());
+ PLANNER_ALL_TODAY_MATCHER.assertMatch(plannerAllTodayRes, plannerAllToday);
+ }
+
+ protected void checkPlannerAllTodayByPlanner(Planner planner) {
+ PlannerAllToday plannerAllToday = PlannerAllTodayBuilder.builder().append(planner).build();
+ PlannerAllToday plannerAllTodayRes = plannerAllTodayImdg.getSingleObjectBySQL(String.format("companyId = %s", planner.getCompanyId()));
+ plannerAllToday.setId(plannerAllTodayRes.getId());
+ PLANNER_ALL_TODAY_MATCHER.assertMatch(plannerAllTodayRes, plannerAllToday);
+ }
+
+ protected ClearingCalendar getClearingCalendar() {
+ ClearingCalendar clearingCalendar = new ClearingCalendar();
+ clearingCalendar.setId(currentID.getAndIncrement());
+ clearingCalendar.setClearingDate(LocalDate.now());
+ clearingCalendar.setCompanyId(deleteCompanyId);
+ clearingCalendar.setDayStatus(DayStatus.Workday.getKey());
+ return clearingCalendar;
+ }
+}
diff --git a/clearing-parent/scheduler-service/src/test/java/ru/spcex/clearing/scheduler/config/SchedulerTestConfig.java b/clearing-parent/scheduler-service/src/test/java/ru/spcex/clearing/scheduler/config/SchedulerTestConfig.java
new file mode 100644
index 000000000..916d3992d
--- /dev/null
+++ b/clearing-parent/scheduler-service/src/test/java/ru/spcex/clearing/scheduler/config/SchedulerTestConfig.java
@@ -0,0 +1,19 @@
+package ru.spcex.clearing.scheduler.config;
+
+import org.springframework.context.annotation.Bean;
+import org.springframework.context.annotation.Configuration;
+import org.springframework.scheduling.TaskScheduler;
+import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler;
+
+@Configuration
+public class SchedulerTestConfig {
+
+ @Bean
+ public TaskScheduler threadPoolTaskScheduler(){
+ ThreadPoolTaskScheduler threadPoolTaskScheduler
+ = new ThreadPoolTaskScheduler();
+ threadPoolTaskScheduler.setPoolSize(5);
+ threadPoolTaskScheduler.setThreadNamePrefix("ThreadPoolTaskScheduler");
+ return threadPoolTaskScheduler;
+ }
+}
diff --git a/clearing-parent/scheduler-service/src/test/java/ru/spcex/clearing/scheduler/service/ClearingCalendarServiceTest.java b/clearing-parent/scheduler-service/src/test/java/ru/spcex/clearing/scheduler/service/ClearingCalendarServiceTest.java
new file mode 100644
index 000000000..fbb7d953b
--- /dev/null
+++ b/clearing-parent/scheduler-service/src/test/java/ru/spcex/clearing/scheduler/service/ClearingCalendarServiceTest.java
@@ -0,0 +1,197 @@
+package ru.spcex.clearing.scheduler.service;
+
+import org.apache.kafka.clients.consumer.MockConsumer;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.springframework.beans.factory.annotation.Autowired;
+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.platform.messaging.domain.BaseRequest;
+import ru.spcex.clearing.platform.messaging.domain.Consts;
+import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteRequest;
+import ru.spcex.clearing.platform.messaging.domain.cud.schedule.ClearingCalendarNewRequest;
+import ru.spcex.clearing.platform.messaging.domain.cud.schedule.ClearingCalendarUpdateRequest;
+import ru.spcex.clearing.scheduler.AbstractServiceTest;
+import ru.spcex.clearing.scheduler.PlannerAllTodayBuilder;
+import ru.spcex.clearing.test.MatcherFactory;
+import ru.spcex.platform.enumeration.DayStatus;
+import ru.spcex.platform.enumeration.Status;
+
+import javax.annotation.PostConstruct;
+import java.time.LocalDate;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static ru.spcex.clearing.test.MatcherFactory.usingIgnoringFieldsComparator;
+import static ru.spcex.clearing.test.TestUtils.*;
+import static ru.spcex.clearing.test.config.ImdgTestConfig.currentID;
+
+class ClearingCalendarServiceTest extends AbstractServiceTest {
+ protected static final MatcherFactory.Matcher CLEARING_CALENDAR_MATCHER = usingIgnoringFieldsComparator("created", "updated");
+ private static final int PARTITION = 0;
+ private static final String TOPIC_CLEARING_CALENDAR_NEW = Consts.DESTINATION_CLEARING_CALENDAR_NEW;
+ private static final String TOPIC_CLEARING_CALENDAR_UPDATE = Consts.DESTINATION_CLEARING_CALENDAR_UPDATE;
+ private static final String TOPIC_CLEARING_CALENDAR_DELETE = Consts.DESTINATION_CLEARING_CALENDAR_DELETE;
+ private static final Long ID = currentID.getAndIncrement();
+
+ @Autowired
+ ClearingCalendarService clearingCalendarService;
+
+ @PostConstruct
+ public void init() {
+ super.init();
+ }
+
+ @BeforeEach
+ public void prepare() {
+ clearAllInImdg(plannerTemplateImdg);
+ clearAllInImdg(plannerAllTodayImdg);
+ }
+
+ /**
+ * {@link ClearingCalendarService#newClearingCalendar(BaseRequest)}(BaseRequest)}
+ * Тест проверяет создание сущности {@link ClearingCalendar} в Hazelcast при передаче из Apache Kafka.
+ * Входной запрос {@link ClearingCalendarNewRequest}:
+ */
+ @Test
+ public void newClearingCalendar() {
+ //ARRANGE
+ ClearingCalendar clearingCalendar = new ClearingCalendar();
+ clearingCalendar.setClearingDate(LocalDate.now());
+ clearingCalendar.setCompanyId(newCompanyId);
+ clearingCalendar.setDayStatus(DayStatus.Workday.getKey());
+ PlannerTemplate plannerTemplate = getPlannerTemplate(newCompanyId);
+ plannerTemplateImdg.insert(plannerTemplate);
+
+ ClearingCalendarNewRequest clearingCalendarNewRequest = new ClearingCalendarNewRequest();
+ clearingCalendarNewRequest.setClearingDate(clearingCalendar.getClearingDate());
+ clearingCalendarNewRequest.setCompanyId(clearingCalendar.getCompanyId());
+ clearingCalendarNewRequest.setDayStatus(clearingCalendar.getDayStatus());
+
+ //ACT
+ String jsonString = getJsonStringForNew(clearingCalendarNewRequest, ID);
+ addRecordToKafka((MockConsumer) clearingCalendarService.getConsumer(), TOPIC_CLEARING_CALENDAR_NEW, PARTITION, 0, jsonString);
+
+ //ASSERT
+ waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord);
+
+ ClearingCalendar clearingCalendarRes = clearingCalendarImdg.getSingleObjectBySQL(String.format("companyId = %s", clearingCalendar.getCompanyId()));
+ clearingCalendar.setId(clearingCalendarRes.getId());
+ CLEARING_CALENDAR_MATCHER.assertMatch(clearingCalendarRes, clearingCalendar);
+ checkPlannerAllTodayByPlannerTemplate(plannerTemplate);
+ }
+
+ /**
+ * {@link ClearingCalendarService#updateClearingCalendar(BaseRequest)}(BaseRequest)}
+ * Тест проверяет создание сущности {@link ClearingCalendar} в Hazelcast при передаче из Apache Kafka.
+ * Входной запрос {@link ClearingCalendarUpdateRequest}:
+ */
+ @Test
+ public void updateClearingCalendar() {
+ //ARRANGE
+ ClearingCalendar clearingCalendar = new ClearingCalendar();
+ clearingCalendar.setId(currentID.getAndIncrement());
+ clearingCalendar.setClearingDate(LocalDate.now());
+ clearingCalendar.setCompanyId(newCompanyId);
+ clearingCalendar.setDayStatus(DayStatus.Workday.getKey());
+ clearingCalendarImdg.insert(clearingCalendar);
+ clearingCalendar.setCompanyId(updateCompanyId);
+
+ PlannerTemplate plannerTemplate = getPlannerTemplate(newCompanyId);
+ plannerTemplateImdg.insert(plannerTemplate);
+
+ ClearingCalendarUpdateRequest clearingCalendarNewRequest = new ClearingCalendarUpdateRequest();
+ clearingCalendarNewRequest.setId(clearingCalendar.getId());
+ clearingCalendarNewRequest.setClearingDate(clearingCalendar.getClearingDate());
+ clearingCalendarNewRequest.setCompanyId(clearingCalendar.getCompanyId());
+ clearingCalendarNewRequest.setDayStatus(clearingCalendar.getDayStatus());
+
+ //ACT
+ String jsonString = getJsonStringForUpdate(clearingCalendarNewRequest, ID);
+ addRecordToKafka((MockConsumer) clearingCalendarService.getConsumer(), TOPIC_CLEARING_CALENDAR_UPDATE, PARTITION, 0, jsonString);
+
+ //ASSERT
+ waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord);
+
+ ClearingCalendar clearingCalendarRes = clearingCalendarImdg.getSingleObjectBySQL(String.format("companyId = %s", clearingCalendar.getCompanyId()));
+ clearingCalendar.setId(clearingCalendarRes.getId());
+ CLEARING_CALENDAR_MATCHER.assertMatch(clearingCalendarRes, clearingCalendar);
+ checkPlannerAllTodayByPlannerTemplate(plannerTemplate);
+ }
+
+
+ /**
+ * {@link ClearingCalendarService#deleteClearingCalendar(BaseRequest)}(BaseRequest)}
+ * Тест проверяет удаление сущности {@link ru.clearing.classes.statics.data.scheduler.ClearingCalendar} в Hazelcast при передаче из Apache Kafka.
+ * Входной запрос {@link CommonDeleteRequest}:
+ */
+ @Test
+ public void deleteClearingCalendar() {
+ //ARRANGE
+ clearAllInImdg(clearingCalendarImdg);
+ ClearingCalendar clearingCalendar = getClearingCalendar();
+ clearingCalendarImdg.insert(clearingCalendar);
+
+ PlannerTemplate plannerTemplate = getPlannerTemplate(newCompanyId);
+ long id = clearingCalendar.getId();
+
+ plannerTemplate.setCompanyId(deleteCompanyId);
+ plannerAllTodayImdg.insert(PlannerAllTodayBuilder.builder().append(plannerTemplate).build());
+
+ CommonDeleteRequest commonDeleteRequest = new CommonDeleteRequest();
+ commonDeleteRequest.setId(id);
+
+ //ACT
+ String jsonString = getJsonStringForDelete(commonDeleteRequest, id);
+ addRecordToKafka((MockConsumer) clearingCalendarService.getConsumer(), TOPIC_CLEARING_CALENDAR_DELETE, PARTITION, 0, jsonString);
+
+ //ASSERT
+ waitingWhenAddedRecordAndCheckIt(id, mockProducer, producerRecord);
+
+ ClearingCalendar plannerTemplateRes = clearingCalendarImdg.getSingleObjectBySQL(String.format("companyId = %s", clearingCalendar.getCompanyId()));
+ assertNull(plannerTemplateRes);
+ PlannerAllToday plannerAllTodayRes = plannerAllTodayImdg.getSingleObjectBySQL(String.format("companyId = %s", plannerTemplate.getCompanyId()));
+ assertNull(plannerAllTodayRes);
+
+ //добавим в clearingCalendarImdg валидный clearingCalendar на случай если тесты plannerTemplate еще не отработали
+ clearingCalendarImdg.insert(clearingCalendar);
+ }
+
+ /**
+ * {@link ClearingCalendarService#deleteFromPlannerAllTodayMap()}}
+ * Тест проверяет удаление сущности {@link ru.clearing.classes.statics.data.scheduler.PlannerAllToday} в Hazelcast.
+ */
+ @Test
+ public void deleteFromPlannerAllTodayMap() {
+ //проверка удаления PlannerAllToday
+ Planner planner = getPlanner(newCompanyId);
+ planner.setId(12L);
+ planner.setTaskStatus(Status.Active.getKey());
+ planner.setCompanyId(newCompanyId);
+ planner.setClearingDate(LocalDate.now());
+ plannerImdg.insert(planner);
+ plannerAllTodayImdg.insert(PlannerAllTodayBuilder.builder().append(planner).build());
+ PlannerTemplate plannerTemplate = new PlannerTemplate();
+ plannerTemplate.setId(13L);
+ plannerTemplate.setCompanyId(updateCompanyId);
+ plannerAllTodayImdg.insert(PlannerAllTodayBuilder.builder().append(plannerTemplate).build());
+
+ //ни чего не должно быть удалено ест валидный clearingCalendar
+ clearingCalendarService.deleteFromPlannerAllTodayMap();
+ assertEquals(2, plannerAllTodayImdg.getAllValues().size());
+ checkPlannerAllTodayByPlanner(planner);
+ checkPlannerAllTodayByPlannerTemplate(plannerTemplate);
+
+ //сейчас удалятся только созданные на основании plannerTemplate
+ clearAllInImdg(clearingCalendarImdg);
+ clearingCalendarService.deleteFromPlannerAllTodayMap();
+
+ checkPlannerAllTodayByPlanner(planner);
+ assertEquals(1, plannerAllTodayImdg.getAllValues().size());
+
+ //добавим в clearingCalendarImdg валидный clearingCalendar на случай если тесты plannerTemplate еще не отработали
+ clearingCalendarImdg.insert(getClearingCalendar());
+ }
+}
\ No newline at end of file
diff --git a/clearing-parent/scheduler-service/src/test/java/ru/spcex/clearing/scheduler/service/LauncherServiceTest.java b/clearing-parent/scheduler-service/src/test/java/ru/spcex/clearing/scheduler/service/LauncherServiceTest.java
new file mode 100644
index 000000000..fc8f4f7a8
--- /dev/null
+++ b/clearing-parent/scheduler-service/src/test/java/ru/spcex/clearing/scheduler/service/LauncherServiceTest.java
@@ -0,0 +1,91 @@
+package ru.spcex.clearing.scheduler.service;
+
+import org.apache.kafka.clients.consumer.MockConsumer;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.springframework.beans.factory.annotation.Autowired;
+import ru.clearing.classes.statics.data.scheduler.Launcher;
+import ru.clearing.classes.statics.data.scheduler.Planner;
+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.LauncherCommandRequest;
+import ru.spcex.clearing.scheduler.AbstractServiceTest;
+import ru.spcex.clearing.test.MatcherFactory;
+
+import javax.annotation.PostConstruct;
+
+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.test.MatcherFactory.usingIgnoringFieldsComparator;
+import static ru.spcex.clearing.test.TestUtils.*;
+import static ru.spcex.clearing.test.config.ImdgTestConfig.currentID;
+import static ru.spcex.clearing.test.config.ImdgTestConfig.defaultAdminId;
+import static ru.spcex.platform.enumeration.Task.accountBlock;
+
+class LauncherServiceTest extends AbstractServiceTest {
+ private final Logger log = LoggerFactory.getLogger(getClass());
+
+ protected static final MatcherFactory.Matcher LAUNCHER_MATCHER = usingIgnoringFieldsComparator("created", "updated");
+ private static final int PARTITION = 0;
+ private static final String TOPIC_LAUNCHER_NEW = Consts.LAUNCHER_NEW;
+ private static final Long ID = currentID.getAndIncrement();
+
+ @Autowired
+ private LauncherService launcherService;
+
+ @PostConstruct
+ public void init() {
+ super.init();
+ }
+
+ @BeforeEach
+ public void prepare() {
+ clearAllInImdg(plannerAllTodayImdg);
+ }
+
+ /**
+ * {@link LauncherService#newLauncher(BaseRequest)}(BaseRequest)}
+ * Тест проверяет создание сущности {@link Planner} в Hazelcast при передаче из Apache Kafka.
+ * Входной запрос {@link LauncherCommandRequest}:
+ */
+ @Test
+ public void newPlanner() {
+ //ARRANGE
+ Launcher launcher = new Launcher();
+ launcher.setSenderId(defaultAdminId);
+ launcher.setTask(accountBlock.getKey());
+
+ LauncherCommandRequest launcherCommandRequest = new LauncherCommandRequest();
+ launcherCommandRequest.setTaskName(launcher.getTask());
+ launcherCommandRequest.setUserId(launcher.getSenderId());
+ launcherCommandRequest.setCompanyId(newCompanyId);
+ launcherCommandRequest.setSecurityId(testSecurityId);
+
+ BaseRequest