From ec0a5edd67ca08ed84f6f3f190809913ee36db5a Mon Sep 17 00:00:00 2001 From: psemenkov Date: Tue, 13 Dec 2022 18:46:42 +0300 Subject: [PATCH 1/4] Adding AccountBalanceServiceTest. --- clearing-parent/balance-service/pom.xml | 37 ++++- .../AccountBalanceValidationRule.java | 4 +- .../balance/config/BalanceImdgTestConfig.java | 75 ++++++++++ .../balance/config/KafkaTestConfig.java | 26 ++++ .../balance/service/AbstractServiceTest.java | 19 +++ .../service/AccountBalanceServiceTest.java | 132 ++++++++++++++++++ .../balance/utils/MatcherFactory.java | 38 +++++ 7 files changed, 327 insertions(+), 4 deletions(-) create mode 100644 clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/config/BalanceImdgTestConfig.java create mode 100644 clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/config/KafkaTestConfig.java create mode 100644 clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/AbstractServiceTest.java create mode 100644 clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/AccountBalanceServiceTest.java create mode 100644 clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/utils/MatcherFactory.java diff --git a/clearing-parent/balance-service/pom.xml b/clearing-parent/balance-service/pom.xml index 1abc9c32c..4414833e5 100644 --- a/clearing-parent/balance-service/pom.xml +++ b/clearing-parent/balance-service/pom.xml @@ -1,6 +1,6 @@ - clearing-parent @@ -40,6 +40,22 @@ ru.spcex.platform platform-enum + + + org.springframework + spring-test + test + + + org.junit.jupiter + junit-jupiter + test + + + org.assertj + assertj-core + test + @@ -67,6 +83,23 @@ ${project.artifactId} + + org.apache.maven.plugins + maven-surefire-plugin + 2.21.0 + + + org.junit.platform + junit-platform-surefire-provider + 1.2.0-M1 + + + org.junit.jupiter + junit-jupiter-engine + 5.2.0-M1 + + + \ No newline at end of file diff --git a/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/validation/AccountBalanceValidationRule.java b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/validation/AccountBalanceValidationRule.java index 7ac689a73..50cac7c02 100644 --- a/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/validation/AccountBalanceValidationRule.java +++ b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/validation/AccountBalanceValidationRule.java @@ -23,7 +23,7 @@ public enum AccountBalanceValidationRule implements IValidationRule companyImdg = context.obtainMap(IMDGDistributedNames.Map_Company, Company.class); Company company = companyImdg.getSingleObjectByID(validatedObject.addresseeId()); if (company == null) { - return of( BalanceError.CompanyNotFound); + return of(BalanceError.CompanyNotFound); } context.storeObject(ValidationStored.Company, company); return empty(); @@ -50,6 +50,6 @@ public enum AccountBalanceValidationRule implements IValidationRule 2) { + pool.setKeepAliveSeconds(60); + pool.setAllowCoreThreadTimeOut(true); + } + pool.setCorePoolSize(maxPoolSz); + pool.setWaitForTasksToCompleteOnShutdown(waitForCompletion); + return pool; + } + + @Bean(name = "taskExecutorHazelcastTestClientInitializer") + public ThreadPoolTaskExecutor taskExecutorHazelcastTestClientInitializer() { + return createThreadPoolTestTaskExecutor(1, true); + } + + @Bean(name = "taskExecutorTestIdGeneratorAwaiter") + public ThreadPoolTaskExecutor taskExecutorTestIdGeneratorAwaiter() { + return createThreadPoolTestTaskExecutor(1, false); + } + + @Autowired + @Bean(name = "hazelcastServiceTest") + public ImdgProvider imdgTestProvider( + @Qualifier("taskExecutorHazelcastTestClientInitializer") ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer, + @Qualifier("taskExecutorTestIdGeneratorAwaiter") ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter, + HazelcastClientParams params) { + Config cfg = new Config(); + cfg.setInstanceName("localhost"); + + NetworkConfig networkConfig = new NetworkConfig(); + JoinConfig joinConfig = new JoinConfig(); + joinConfig.setMulticastConfig(new MulticastConfig().setEnabled(false)); + joinConfig.setTcpIpConfig(new TcpIpConfig().setEnabled(true).setMembers(List.of("127.0.0.1"))); + networkConfig.setJoin(joinConfig); + cfg.setNetworkConfig(networkConfig); + hazelcastInstance = Hazelcast.newHazelcastInstance(cfg); + HazelcastHelper.otcSystem_setStorageState(true, hazelcastInstance); + return new HazelcastService(taskExecutorHazelcastClientInitializer, taskExecutorIdGeneratorAwaiter, params); + } + + @Bean(name = "hazelcastClientParams") + public HazelcastClientParams getHazelcastClientParams() { + HazelcastClientParams params = new HazelcastClientParams(); + params.setLogin("dev"); + params.setPassword("dev-pass"); + params.setClusterMembers("127.0.0.1"); + params.setInstanceName("hzTestClient" + new Random().nextInt()); + params.setNearCacheConfig(new NearCacheConfig()); + return params; + } +} diff --git a/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/config/KafkaTestConfig.java b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/config/KafkaTestConfig.java new file mode 100644 index 000000000..33e62153e --- /dev/null +++ b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/config/KafkaTestConfig.java @@ -0,0 +1,26 @@ +package ru.spcex.clearing.balance.config; + +import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.consumer.MockConsumer; +import org.apache.kafka.clients.consumer.OffsetResetStrategy; +import org.apache.kafka.clients.producer.MockProducer; +import org.apache.kafka.clients.producer.Producer; +import org.springframework.beans.factory.config.ConfigurableBeanFactory; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.Scope; + +@Configuration +public class KafkaTestConfig { + + @Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) + @Bean + public Consumer createTestConsumer() { + return new MockConsumer<>(OffsetResetStrategy.EARLIEST); + } + + @Bean + public Producer createTestProducer() { + return new MockProducer<>(); + } +} diff --git a/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/AbstractServiceTest.java b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/AbstractServiceTest.java new file mode 100644 index 000000000..29723228e --- /dev/null +++ b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/AbstractServiceTest.java @@ -0,0 +1,19 @@ +package ru.spcex.clearing.balance.service; + +import org.junit.jupiter.api.extension.ExtendWith; +import org.springframework.test.context.ContextConfiguration; +import org.springframework.test.context.junit.jupiter.SpringExtension; +import ru.spcex.clearing.balance.config.BalanceImdgTestConfig; +import ru.spcex.clearing.balance.config.KafkaSenderConfig; +import ru.spcex.clearing.balance.config.KafkaTestConfig; +import ru.spcex.clearing.balance.config.ValidationConfig; + +@ExtendWith(SpringExtension.class) +@ContextConfiguration(classes = { + AccountBalanceService.class, + ValidationConfig.class, + BalanceImdgTestConfig.class, + KafkaSenderConfig.class, + KafkaTestConfig.class}) +public abstract class AbstractServiceTest { +} diff --git a/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/AccountBalanceServiceTest.java b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/AccountBalanceServiceTest.java new file mode 100644 index 000000000..31a79fe58 --- /dev/null +++ b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/AccountBalanceServiceTest.java @@ -0,0 +1,132 @@ +package ru.spcex.clearing.balance.service; + +import org.junit.jupiter.api.Test; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; +import ru.clearing.classes.statics.data.account.Account; +import ru.clearing.classes.statics.data.account.AccountBalance; +import ru.clearing.classes.statics.data.company.Company; +import ru.spcex.clearing.balance.errors.BalanceError; +import ru.spcex.clearing.balance.utils.MatcherFactory.Matcher; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.platform.enumeration.AccountType; +import ru.spcex.platform.enumeration.BalanceAccountType; +import ru.spcex.platform.enumeration.Status; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; +import ru.spcex.platform.utils.enumeration.EnumMessage; + +import javax.annotation.PostConstruct; +import java.math.BigDecimal; + +import static ru.spcex.clearing.balance.utils.MatcherFactory.usingIgnoringFieldsComparator; + + +class AccountBalanceServiceTest extends AbstractServiceTest { + + public static final Matcher ACCOUNT_BALANCE_MATCHER = usingIgnoringFieldsComparator("account.created", "account.updated", "account.clearingDate"); + private final Long id = 0L; + private final Long accountIdNew = 0L; + private final Long addresseeIdNew = 0L; + private final BigDecimal amountNew = new BigDecimal(1000); + private final String cashMovementCurrencyCodeNew = "cash"; + private final Long accountIdUpdate = accountIdNew; + private final Long addresseeIdUpdate = addresseeIdNew; + private final BigDecimal amountUpdate = new BigDecimal(100000); + private final String cashMovementCurrencyCodeUpdate = "cash"; + @Autowired + AccountBalanceService accountBalanceService; + @Autowired + @Qualifier("hazelcastServiceTest") + private ImdgProvider hazelcast; + + private Imdg companyImdg; + private Imdg accountImdg; + + @PostConstruct + void init() { + companyImdg = hazelcast.getImdg(IMDGDistributedNames.Map_Company, Company.class); + accountImdg = hazelcast.getImdg(IMDGDistributedNames.Map_Account, Account.class); + } + + @Test + void createAccountBalance() { + hazelcast.waitAvailable(); + Company company = new Company(); + company.setId(addresseeIdNew); + company.setTradingCode("code"); + company.setShortName("ShortName"); + company.setFullName("FullName"); + companyImdg.insert(company); + Account account = new Account(); + account.setId(accountIdNew); + account.setAccount("123456789"); + account.setAccountType(AccountType.Clrn.getKey()); + account.setAccountStatus(Status.Active.getKey()); + accountImdg.insert(account); + + AccountBalance accountBalanceNew = new AccountBalance(); + accountBalanceNew.setId(id); + accountBalanceNew.setCompanyId(addresseeIdNew); + accountBalanceNew.setAccountId(accountIdNew); + accountBalanceNew.setAccountType(account.getAccountType()); + accountBalanceNew.setAccount(account.getAccount()); + accountBalanceNew.setOpenBalanceAmount(amountNew); + accountBalanceNew.setFreeBalanceAmount(amountNew); + accountBalanceNew.setBalanceAmount(amountNew); + accountBalanceNew.setBalanceAccountType(BalanceAccountType.Active.getKey()); + accountBalanceNew.setCurrencyCode(cashMovementCurrencyCodeNew); + accountBalanceNew.setTradingCode(company.getTradingCode()); + accountBalanceNew.setShortName(company.getShortName()); + accountBalanceNew.setFullName(company.getFullName()); + AccountResult predictableNewResult = new AccountResult(accountBalanceNew); + AccountResult resultNew = accountBalanceService.createAccountBalance(addresseeIdNew, accountIdNew, amountNew, cashMovementCurrencyCodeNew); + resultNew.getAccount().setId(id); + + AccountBalance accountBalanceUpdate = new AccountBalance(); + accountBalanceUpdate.setId(id); + accountBalanceUpdate.setCompanyId(addresseeIdUpdate); + accountBalanceUpdate.setAccountId(accountIdUpdate); + accountBalanceUpdate.setAccountType(account.getAccountType()); + accountBalanceUpdate.setAccount(account.getAccount()); + accountBalanceUpdate.setOpenBalanceAmount(amountUpdate); + accountBalanceUpdate.setFreeBalanceAmount(amountUpdate); + accountBalanceUpdate.setBalanceAmount(amountUpdate); + accountBalanceUpdate.setBalanceAccountType(BalanceAccountType.Active.getKey()); + accountBalanceUpdate.setCurrencyCode(cashMovementCurrencyCodeUpdate); + accountBalanceUpdate.setTradingCode(company.getTradingCode()); + accountBalanceUpdate.setShortName(company.getShortName()); + accountBalanceUpdate.setFullName(company.getFullName()); + AccountResult resultUpdate = accountBalanceService.createAccountBalance(addresseeIdUpdate, accountIdUpdate, amountUpdate, cashMovementCurrencyCodeUpdate); + resultUpdate.getAccount().setId(id); + AccountResult predictableUpdateResult = new AccountResult(accountBalanceUpdate); + + ACCOUNT_BALANCE_MATCHER.assertMatch(resultNew, predictableNewResult); + ACCOUNT_BALANCE_MATCHER.assertMatch(resultUpdate, predictableUpdateResult); + } + + @Test + void validatedCreateAccountBalance() { + AccountResult predictableResult; + hazelcast.waitAvailable(); + predictableResult = new AccountResult(new EnumMessage(BalanceError.CompanyNotFound)); + + AccountResult result = accountBalanceService.createAccountBalance(addresseeIdNew, accountIdNew, amountNew, cashMovementCurrencyCodeNew); + ACCOUNT_BALANCE_MATCHER.assertMatch(result, predictableResult); + + Company company = new Company(); + company.setId(addresseeIdNew); + companyImdg.insert(company); + + Account account = new Account(); +// account.setId(accountIdNew); + account.setAccountType(AccountType.Clrn.getKey()); + account.setAccountStatus(Status.Active.getKey()); + accountImdg.insert(account); + predictableResult = new AccountResult(new EnumMessage(BalanceError.AccountNotPresent)); + + result = accountBalanceService.createAccountBalance(addresseeIdNew, accountIdNew, amountNew, cashMovementCurrencyCodeNew); + ACCOUNT_BALANCE_MATCHER.assertMatch(result, predictableResult); + + } +} \ No newline at end of file diff --git a/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/utils/MatcherFactory.java b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/utils/MatcherFactory.java new file mode 100644 index 000000000..a18ccb4f0 --- /dev/null +++ b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/utils/MatcherFactory.java @@ -0,0 +1,38 @@ +package ru.spcex.clearing.balance.utils; + +import java.util.Arrays; + +import static org.assertj.core.api.Assertions.assertThat; + +/** + * Factory for creating test matchers. + *

+ * Comparing actual and expected objects via AssertJ + */ +public class MatcherFactory { + + public static Matcher usingIgnoringFieldsComparator(String... fieldsToIgnore) { + return new Matcher<>(fieldsToIgnore); + } + + public static class Matcher { + private final String[] fieldsToIgnore; + + private Matcher(String... fieldsToIgnore) { + this.fieldsToIgnore = fieldsToIgnore; + } + + public void assertMatch(T actual, T expected) { + assertThat(actual).usingRecursiveComparison().ignoringFields(fieldsToIgnore).isEqualTo(expected); + } + + @SafeVarargs + public final void assertMatch(Iterable actual, T... expected) { + assertMatch(actual, Arrays.asList(expected)); + } + + public void assertMatch(Iterable actual, Iterable expected) { + assertThat(actual).usingRecursiveFieldByFieldElementComparatorIgnoringFields(fieldsToIgnore).isEqualTo(expected); + } + } +} From 9830cc88059eedc83dd35db3f9063fc819a95ece Mon Sep 17 00:00:00 2001 From: aalehin Date: Tue, 13 Dec 2022 19:56:41 +0300 Subject: [PATCH 2/4] task manager... --- .../scheduler/enums/IEnumWithLongValue.java | 81 +++++++ .../scheduler/enums/TaskStatuses.java | 31 +++ .../spcex/clearing/scheduler/enums/Tasks.java | 20 ++ .../scheduler/service/TaskManager.java | 219 ++++++++++++++++++ 4 files changed, 351 insertions(+) create mode 100644 clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/enums/IEnumWithLongValue.java create mode 100644 clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/enums/TaskStatuses.java create mode 100644 clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/enums/Tasks.java create mode 100644 clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/TaskManager.java diff --git a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/enums/IEnumWithLongValue.java b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/enums/IEnumWithLongValue.java new file mode 100644 index 000000000..3e46ccc59 --- /dev/null +++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/enums/IEnumWithLongValue.java @@ -0,0 +1,81 @@ +package ru.spcex.clearing.scheduler.enums; + +import java.io.Serializable; +import java.util.Objects; + +/** + * Удобный интерфейс для enum при использовании TextErrorService. + */ +public interface IEnumWithLongValue extends Serializable { + + + /** + * Проверяет что среди данного набора Enum, присутствует элемент с данным id + * id может быть null + */ + static & IEnumWithLongValue> boolean contains(Long id, T... enumSet) { + for (T e : enumSet) { + if (Objects.equals(id, e.getId())) + return true; + } + return false; + } + + static & IEnumWithLongValue> boolean contains(T e, T... enumSet) { + for (T enumEl : enumSet) { + if (enumEl.equals(e)) + return true; + } + return false; + } + + static & IEnumWithLongValue> Long[] toLongArray(T... enumSet) { + Long[] enumId = new Long[enumSet.length]; + for (int i = 0; i < enumSet.length; i++) + enumId[i] = enumSet[i].getId(); + return enumId; + } + + /** + * Проверяет что среди всех Enum данного класса, присутствует элемент с данным id + * id может быть null + */ + static & IEnumWithLongValue> boolean contains(Class enumClass, Long id) { + for (T e : enumClass.getEnumConstants()) { + if (Objects.equals(id, e.getId())) + return true; + } + return false; + } + + + /** + * Возвращает Enum по id, если в заданном классе такой определен + * id может быть null + * + * @return Enum если нашел, иначе null + */ + static & IEnumWithLongValue> T getEnumById(Class enumClass, Long id) { + for (T e : enumClass.getEnumConstants()) { + if (Objects.equals(id, e.getId())) + return e; + } + return null; + } + + static Long getIdOrNull(IEnumWithLongValue enumVal) { + return enumVal == null ? null : enumVal.getId(); + } + + + /** + * errorCode + **/ + Long getId(); + + default boolean equalsById(Long id) { + return id != null && getId().equals(id); + } + + String elementName(); +} diff --git a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/enums/TaskStatuses.java b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/enums/TaskStatuses.java new file mode 100644 index 000000000..ab2709ba8 --- /dev/null +++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/enums/TaskStatuses.java @@ -0,0 +1,31 @@ +package ru.spcex.clearing.scheduler.enums; + +import java.util.Objects; + +public enum TaskStatuses { + ACTIVE("ACTV"), + BLOCKED("CNCL"), + CANCEL("BLKD"); + + private final String name; + + TaskStatuses(String name) { + this.name = name; + } + + static TaskStatuses getEnumById(String name) { + for (TaskStatuses e : TaskStatuses.values()) { + if (Objects.equals(name, e.name())) + return e; + } + return null; + } + + public String getName() { + return this.name; + } + + public String elementName() { + return this.name(); + } +} diff --git a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/enums/Tasks.java b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/enums/Tasks.java new file mode 100644 index 000000000..7a5ec165b --- /dev/null +++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/enums/Tasks.java @@ -0,0 +1,20 @@ +package ru.spcex.clearing.scheduler.enums; + +public enum Tasks implements IEnumWithLongValue { + SOME_STATUS(1L);//TODO CLEARIFY + + private final Long id; + + Tasks(Long id) { + this.id = id; + } + + public Long getId() { + return this.id; + } + + public String elementName() { + return this.name(); + } + +} 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 new file mode 100644 index 000000000..7e226ffd5 --- /dev/null +++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/TaskManager.java @@ -0,0 +1,219 @@ +package ru.spcex.clearing.scheduler.service; + +import com.hazelcast.core.EntryEvent; +import com.hazelcast.core.HazelcastInstance; +import com.hazelcast.core.IMap; +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.scheduling.TaskScheduler; +import ru.clearing.classes.statics.data.scheduler.PlannerAllToday; +import ru.spcex.clearing.scheduler.enums.Tasks; + +import java.time.LocalDate; +import java.time.LocalDateTime; +import java.time.LocalTime; +import java.time.ZoneId; +import java.util.Date; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ScheduledFuture; + +import static ru.spcex.clearing.imdg.IMDGDistributedNames.Map_PlannerAllToday; + +/** + * Планировщик задач, расписание берёт из Hazelcast map. + *

+ */ +public abstract class TaskManager implements EntryAddedListener, + EntryUpdatedListener, EntryRemovedListener, + InitializingBean { + private static final Logger log = LoggerFactory.getLogger(TaskManager.class); + + protected HazelcastInstance hazelcastInstance; + protected IMap plannerAllTodayMapStore; + protected TaskScheduler taskScheduler; + + protected ConcurrentHashMap scheduledJobs; + + protected TaskManager(TaskScheduler taskScheduler, HazelcastInstance hazelcastInstance) { + this.taskScheduler = taskScheduler; + this.hazelcastInstance = hazelcastInstance; + } + + private static LocalDateTime dateOldTypeConvert(Date oldDate) { + return LocalDateTime.ofInstant(oldDate.toInstant(), ZoneId.systemDefault()); + } + + private static LocalDate dateTypeConvert(Date oldDate) { + return dateOldTypeConvert(oldDate).toLocalDate(); + } + + private static LocalTime timeTypeConvert(Date oldDate) { + return dateOldTypeConvert(oldDate).toLocalTime(); + } + + @Override + public void afterPropertiesSet() { + plannerAllTodayMapStore = hazelcastInstance.getMap(Map_PlannerAllToday); + scheduledJobs = new ConcurrentHashMap<>(); + updateScheduler(); + } + + // --- Слушатели Hazelcast Map --- + @Override + public void entryAdded(EntryEvent event) { + PlannerAllToday task = event.getValue(); + + + //todo define how to get task status Tasks taskType = getEnumById(Tasks.class, task.getParentId()); +// if (taskType == null) { +// log.warn("Task skipped, {} task type not recognized", task.getTaskId()); +// return; +// } + processTask(task); + } + + @Override + public void entryUpdated(EntryEvent event) { + PlannerAllToday task = event.getValue(); + PlannerAllToday oldTask = event.getOldValue(); + + LocalTime oldTime = oldTask.getTaskTime(); +// todo define TaskStatuses oldStatus + //если таск относится к другому обработчику, пропускаем +// if (!getTaskId().equalsById(task.getTaskId()) && !getTaskId().equalsById(oldTask.getTaskId())) { +// return; +// } +// if (!Objects.equals(task.getTaskId(), oldTask.getTaskId())) { +// throw new IllegalStateException("changed taskId for SchedulerAllToday in core"); +// } +// if (ACTIVE.equalsById(oldTask.getTaskStatusId())) { //&& ACTIVE.equalsById(task.getTaskStatusId()) +// if (!removeTask(oldTime)) { +// log.debug("cannot cancel task with type {}, time {}", getTaskId().name(), oldTime.toString()); +// } else { +// log.debug("task with type {}, time {} execution cancelled, adding altered task...", getTaskId().name(), oldTime.toString()); +// } +// processTask(task); +// } else if (CANCEL.equalsById(oldTask.getTaskStatusId())) { //&& TaskStatuses.CANCEL.equalsById(task.getTaskStatusId()) +// restorePreviouslyRemovedTask(oldTime, oldTask); +// processTask(task); +// } else if (BLOCKED.equalsById(oldTask.getTaskStatusId())) { +// processTask(task); +// } + } + + @Override + public void entryRemoved(EntryEvent event) { + PlannerAllToday taskToRemove = event.getOldValue(); + if (taskToRemove == null) { + log.debug("no task in removed event"); + return; + } + //LocalTime removedTaskTime = timeTypeConvert(taskToRemove.getTaskTime()); + // + //if (ACTIVE.equalsById(taskToRemove.getTaskStatusId())) { + // if (removeTask(removedTaskTime)) + // log.debug("task successfully canceled"); + //} else if (CANCEL.equalsById(taskToRemove.getTaskStatusId())) { + // restorePreviouslyRemovedTask(removedTaskTime, taskToRemove); + //} else if (BLOCKED.equalsById(taskToRemove.getTaskStatusId())) { + // log.debug("BLOCKED task removed; do nothing"); + //} + } + + private void restorePreviouslyRemovedTask(LocalTime oldTime, PlannerAllToday oldTask) { + ScheduledFuture cancelledFuture = scheduledJobs.get(oldTime); + if (cancelledFuture != null && cancelledFuture.isCancelled()) { + PlannerAllToday schedulerAllToday = new PlannerAllToday(); + // schedulerAllToday.setTaskId(getTaskId().getId()); + // schedulerAllToday.setTaskStatusId(ACTIVE.getId()); + schedulerAllToday.setTaskTime(oldTask.getTaskTime()); + log.debug("CANCEL task updated/removed; restoring previously cancelled task"); + processTask(schedulerAllToday); + } + } + + // --- Планирование задач --- + protected void updateScheduler() { + //Collection schedulerAllTodays = plannerAllTodayMapStore.values(); + //Collection sortedSchedulers = + // schedulerAllTodays.stream().sorted((o1, o2) -> (ACTIVE.equalsById(o1.getTaskStatusId()) && CANCEL.equalsById(o2.getTaskStatusId())) ? -1 : 0) + // .collect(Collectors.toCollection(ArrayList::new)); + //for (PlannerAllToday schedulerAllToday : sortedSchedulers) { + // processTask(schedulerAllToday); + //} + } + + // --- Работа с задачами --- + private void processTask(PlannerAllToday task) { + //LocalTime taskTime = task.getTaskTime(); + //Tasks taskType = getEnumById(Tasks.class, task.getTaskId()); + //TaskStatuses taskStatus = getEnumById(TaskStatuses.class, task.getTaskStatusId()); + //if (taskStatus == null) throw new IllegalStateException("task status from core can't be null"); + //if (!getTaskId().equals(taskType)) { + // log.debug("Task skipped - taskTime {}, is not of acceptable type {}", taskTime.toString(), taskType != null ? taskType.name() : ""); + // return; + //} + //if (taskTime.isBefore(LocalTime.now())) { + // log.debug("Task skipped - taskTime {} is before now, taskType {} ok", taskTime, getTaskId().toString()); + // return; + //} + //if (taskStatus.equals(BLOCKED)) { + // log.debug("Task skipped - timeTime {} with status {}", taskTime.toString(), BLOCKED.toString()); + // return; + //} + //{ + // ScheduledFuture future = scheduledJobs.get(taskTime); + // if (future != null) { + // if (TaskStatuses.CANCEL.equals(taskStatus)) { + // //случай когда пришел cancel, пытаемся отменить зарегистрированный ранее таск + // future.cancel(false); + // log.debug("cancelling task time {}", taskTime); + // return; + // } else if (TaskStatuses.ACTIVE.equals(taskStatus)) { + // //случай когда пришел активный таск, и уже был на это время неотмененный + // if (!future.isCancelled()) { + // log.debug("such task time {} has already been registered", taskTime); + // return; + // } + // } + // } + //} + //if (TaskStatuses.CANCEL.equals(taskStatus)) { + // log.debug("task type {} time {} with status CANCEL - no tasks to cancel found", taskType.name(), taskTime); + // return; + //} + ////пришел активный таск + //log.debug("adding task type {}, time {}", taskType.name(), taskTime); + //ScheduledFuture future = taskScheduler.schedule(() -> doJob(task), LocalDateTime.of(LocalDate.now(), taskTime).atZone(ZoneId.systemDefault()).toInstant()); + //ScheduledFuture oldFuture = scheduledJobs.put(taskTime, future); + //if (oldFuture != null && !oldFuture.isCancelled()) { //for synchronization, never + // log.warn("tasks were added simultaneously, cancel former one"); + // boolean success = oldFuture.cancel(false); + // log.warn("cancelling task " + (success ? "success" : "fail")); + //} + } + + private boolean removeTask(LocalTime taskTime) { + ScheduledFuture future = scheduledJobs.remove(taskTime); + if (future == null) { + log.debug("can't cancel task type {}, time {}, not found", getTaskId().name(), taskTime); + return false; + } + return future.cancel(false); + } + + // --- Реализация выполнения задач --- + + /** + * Рабочий тип запланированных задач. + * + * @return идентификатор для фильтра типов планировщика задачь. Планировать задачи только этого типа. + */ + protected abstract Tasks getTaskId(); + + protected abstract void doJob(PlannerAllToday taskInfo); +} From b225e0d3d98462bf55fe8a3c4babc65c927b2d52 Mon Sep 17 00:00:00 2001 From: aalehin Date: Wed, 14 Dec 2022 10:41:37 +0300 Subject: [PATCH 3/4] =?UTF-8?q?http://jira.mfd.msk:8088/browse/CLS-37=20--?= =?UTF-8?q?-=20=D0=BF=D0=B5=D1=80=D0=B5=D0=BD=D0=B5=D1=81=20task=20manager?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../scheduler/enums/TaskStatuses.java | 6 +- .../spcex/clearing/scheduler/enums/Tasks.java | 28 ++- .../scheduler/service/TaskManager.java | 191 +++++++++--------- 3 files changed, 124 insertions(+), 101 deletions(-) diff --git a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/enums/TaskStatuses.java b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/enums/TaskStatuses.java index ab2709ba8..b4f82ea38 100644 --- a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/enums/TaskStatuses.java +++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/enums/TaskStatuses.java @@ -13,7 +13,7 @@ public enum TaskStatuses { this.name = name; } - static TaskStatuses getEnumById(String name) { + public static TaskStatuses getEnumByName(String name) { for (TaskStatuses e : TaskStatuses.values()) { if (Objects.equals(name, e.name())) return e; @@ -21,6 +21,10 @@ public enum TaskStatuses { return null; } + public Boolean equalsByName(String name) { + return this.name.equalsIgnoreCase(name); + } + public String getName() { return this.name; } diff --git a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/enums/Tasks.java b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/enums/Tasks.java index 7a5ec165b..be73dfe08 100644 --- a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/enums/Tasks.java +++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/enums/Tasks.java @@ -1,16 +1,30 @@ package ru.spcex.clearing.scheduler.enums; -public enum Tasks implements IEnumWithLongValue { - SOME_STATUS(1L);//TODO CLEARIFY +import java.util.Objects; - private final Long id; +public enum Tasks { + SOME_STATUS("STATUS"); - Tasks(Long id) { - this.id = id; + private String name; + + Tasks(String name) { + this.name = name; } - public Long getId() { - return this.id; + public static Tasks getEnumByName(String name) { + for (Tasks e : Tasks.values()) { + if (Objects.equals(name, e.name())) + return e; + } + return null; + } + + public Boolean equalsByName(String name) { + return this.name.equalsIgnoreCase(name); + } + + public String getName() { + return this.name; } public String elementName() { 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 7e226ffd5..4bfd70cdb 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 @@ -11,17 +11,23 @@ import org.slf4j.LoggerFactory; import org.springframework.beans.factory.InitializingBean; import org.springframework.scheduling.TaskScheduler; import ru.clearing.classes.statics.data.scheduler.PlannerAllToday; +import ru.spcex.clearing.scheduler.enums.TaskStatuses; import ru.spcex.clearing.scheduler.enums.Tasks; import java.time.LocalDate; import java.time.LocalDateTime; import java.time.LocalTime; import java.time.ZoneId; +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.stream.Collectors; import static ru.spcex.clearing.imdg.IMDGDistributedNames.Map_PlannerAllToday; +import static ru.spcex.clearing.scheduler.enums.TaskStatuses.*; /** * Планировщик задач, расписание берёт из Hazelcast map. @@ -67,12 +73,11 @@ public abstract class TaskManager implements EntryAddedListener event) { PlannerAllToday task = event.getValue(); - - //todo define how to get task status Tasks taskType = getEnumById(Tasks.class, task.getParentId()); -// if (taskType == null) { -// log.warn("Task skipped, {} task type not recognized", task.getTaskId()); -// return; -// } + Tasks taskType = Tasks.getEnumByName(task.getTask()); + if (taskType == null) { + log.warn("Task skipped, {} task type not recognized", task.getTask()); + return; + } processTask(task); } @@ -82,27 +87,27 @@ public abstract class TaskManager implements EntryAddedListener schedulerAllTodays = plannerAllTodayMapStore.values(); - //Collection sortedSchedulers = - // schedulerAllTodays.stream().sorted((o1, o2) -> (ACTIVE.equalsById(o1.getTaskStatusId()) && CANCEL.equalsById(o2.getTaskStatusId())) ? -1 : 0) - // .collect(Collectors.toCollection(ArrayList::new)); - //for (PlannerAllToday schedulerAllToday : sortedSchedulers) { - // processTask(schedulerAllToday); - //} + Collection schedulerAllTodays = plannerAllTodayMapStore.values(); + Collection sortedSchedulers = + schedulerAllTodays.stream().sorted((o1, o2) -> (ACTIVE.equalsByName(o1.getTaskStatus()) && CANCEL.equalsByName(o2.getTaskStatus())) ? -1 : 0) + .collect(Collectors.toCollection(ArrayList::new)); + for (PlannerAllToday schedulerAllToday : sortedSchedulers) { + processTask(schedulerAllToday); + } } // --- Работа с задачами --- private void processTask(PlannerAllToday task) { - //LocalTime taskTime = task.getTaskTime(); - //Tasks taskType = getEnumById(Tasks.class, task.getTaskId()); - //TaskStatuses taskStatus = getEnumById(TaskStatuses.class, task.getTaskStatusId()); - //if (taskStatus == null) throw new IllegalStateException("task status from core can't be null"); - //if (!getTaskId().equals(taskType)) { - // log.debug("Task skipped - taskTime {}, is not of acceptable type {}", taskTime.toString(), taskType != null ? taskType.name() : ""); - // return; - //} - //if (taskTime.isBefore(LocalTime.now())) { - // log.debug("Task skipped - taskTime {} is before now, taskType {} ok", taskTime, getTaskId().toString()); - // return; - //} - //if (taskStatus.equals(BLOCKED)) { - // log.debug("Task skipped - timeTime {} with status {}", taskTime.toString(), BLOCKED.toString()); - // return; - //} - //{ - // ScheduledFuture future = scheduledJobs.get(taskTime); - // if (future != null) { - // if (TaskStatuses.CANCEL.equals(taskStatus)) { - // //случай когда пришел cancel, пытаемся отменить зарегистрированный ранее таск - // future.cancel(false); - // log.debug("cancelling task time {}", taskTime); - // return; - // } else if (TaskStatuses.ACTIVE.equals(taskStatus)) { - // //случай когда пришел активный таск, и уже был на это время неотмененный - // if (!future.isCancelled()) { - // log.debug("such task time {} has already been registered", taskTime); - // return; - // } - // } - // } - //} - //if (TaskStatuses.CANCEL.equals(taskStatus)) { - // log.debug("task type {} time {} with status CANCEL - no tasks to cancel found", taskType.name(), taskTime); - // return; - //} - ////пришел активный таск - //log.debug("adding task type {}, time {}", taskType.name(), taskTime); - //ScheduledFuture future = taskScheduler.schedule(() -> doJob(task), LocalDateTime.of(LocalDate.now(), taskTime).atZone(ZoneId.systemDefault()).toInstant()); - //ScheduledFuture oldFuture = scheduledJobs.put(taskTime, future); - //if (oldFuture != null && !oldFuture.isCancelled()) { //for synchronization, never - // log.warn("tasks were added simultaneously, cancel former one"); - // boolean success = oldFuture.cancel(false); - // log.warn("cancelling task " + (success ? "success" : "fail")); - //} + LocalTime taskTime = task.getTaskTime(); + Tasks taskType = Tasks.getEnumByName(task.getTask()); + TaskStatuses taskStatus = TaskStatuses.getEnumByName(task.getTaskStatus()); + if (taskStatus == null) throw new IllegalStateException("task status from core can't be null"); + if (!getTask().equals(taskType)) { + log.debug("Task skipped - taskTime {}, is not of acceptable type {}", taskTime.toString(), taskType != null ? taskType.name() : ""); + return; + } + if (taskTime.isBefore(LocalTime.now())) { + log.debug("Task skipped - taskTime {} is before now, taskType {} ok", taskTime, getTask().toString()); + return; + } + if (taskStatus.equals(BLOCKED)) { + log.debug("Task skipped - timeTime {} with status {}", taskTime.toString(), BLOCKED.toString()); + return; + } + { + ScheduledFuture future = scheduledJobs.get(taskTime); + if (future != null) { + if (TaskStatuses.CANCEL.equals(taskStatus)) { + //случай когда пришел cancel, пытаемся отменить зарегистрированный ранее таск + future.cancel(false); + log.debug("cancelling task time {}", taskTime); + return; + } else if (TaskStatuses.ACTIVE.equals(taskStatus)) { + //случай когда пришел активный таск, и уже был на это время неотмененный + if (!future.isCancelled()) { + log.debug("such task time {} has already been registered", taskTime); + return; + } + } + } + } + if (TaskStatuses.CANCEL.equals(taskStatus)) { + log.debug("task type {} time {} with status CANCEL - no tasks to cancel found", taskType.name(), taskTime); + return; + } + //пришел активный таск + log.debug("adding task type {}, time {}", taskType.name(), taskTime); + ScheduledFuture future = taskScheduler.schedule(() -> doJob(task), LocalDateTime.of(LocalDate.now(), taskTime).atZone(ZoneId.systemDefault()).toInstant()); + ScheduledFuture oldFuture = scheduledJobs.put(taskTime, future); + if (oldFuture != null && !oldFuture.isCancelled()) { //for synchronization, never + log.warn("tasks were added simultaneously, cancel former one"); + boolean success = oldFuture.cancel(false); + log.warn("cancelling task " + (success ? "success" : "fail")); + } } private boolean removeTask(LocalTime taskTime) { ScheduledFuture future = scheduledJobs.remove(taskTime); if (future == null) { - log.debug("can't cancel task type {}, time {}, not found", getTaskId().name(), taskTime); + log.debug("can't cancel task type {}, time {}, not found", getTask().name(), taskTime); return false; } return future.cancel(false); @@ -213,7 +218,7 @@ public abstract class TaskManager implements EntryAddedListener Date: Wed, 14 Dec 2022 10:45:22 +0300 Subject: [PATCH 4/4] =?UTF-8?q?http://jira.mfd.msk:8088/browse/CLS-37=20--?= =?UTF-8?q?-=20=D0=B4=D0=BE=D0=B1=D0=B0=D0=B2=D0=B8=D0=BB=20=D0=B7=D0=BD?= =?UTF-8?q?=D0=B0=D1=87=D0=B5=D0=BD=D0=B8=D1=8F=20=D0=B7=D0=B0=D0=B4=D0=B0?= =?UTF-8?q?=D1=87?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../ru/spcex/clearing/scheduler/enums/Tasks.java | 13 ++++++++++++- 1 file changed, 12 insertions(+), 1 deletion(-) diff --git a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/enums/Tasks.java b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/enums/Tasks.java index be73dfe08..b92b42108 100644 --- a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/enums/Tasks.java +++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/enums/Tasks.java @@ -3,7 +3,18 @@ package ru.spcex.clearing.scheduler.enums; import java.util.Objects; public enum Tasks { - SOME_STATUS("STATUS"); + GBAL("GBAL"),//Зачисление остатков + ABLK("ABLK"),//Блокировка счета + GALB("GALB"),//Запрос остатков по всем счетам + ADBL("ADBL"),//Дозачисление/списание остатков + CORD("CORD"),//Формирование сводного платежного поручения + CORC("CORC"),//Получение подтверждения переводов + GTRD("GTRD"),//Получение сделок из Торговой системы + GVER("GVER"),// Запуск сверки + GBLD("GBLD"),// Поступление средств + SCLR("SCLR"),// Запуск клиринговой сессии + SPRC("SPRC"),// Запуск преклиринга + SPOC("SPOC");// Запуск постклиринга private String name;